From cac8e5ae8e917c45ab5107354f876aa16e092674 Mon Sep 17 00:00:00 2001 From: Leandro Moraes Date: Thu, 23 Jul 2026 09:53:26 -0300 Subject: [PATCH 1/4] =?UTF-8?q?fix(outbox):=20corrige=20self-invocation=20?= =?UTF-8?q?que=20anulava=20@Transactional=20+=20reclama=20eventos=20IN=5FF?= =?UTF-8?q?LIGHT=20=C3=B3rf=C3=A3os?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Extrai claimBatch()/markPublished()/revertToPending() de OutboxRelay pra OutboxClaimCoordinator, um @Component novo. OutboxRelay chamava esses métodos via this.*, o que faz o proxy AOP do Spring nunca interceptar a chamada e deixa @Transactional como no-op silencioso. Agora OutboxRelay recebe o coordinator por injeção de dependência e chama através dele, passando pelo proxy de verdade. Adiciona reaper de eventos IN_FLIGHT travados: claimBatch() primeiro reclama de volta pra PENDING qualquer evento IN_FLIGHT com claimed_at mais velho que 2 minutos, antes de buscar o próximo lote PENDING. Nova coluna claimed_at (migration V4) marca quando um evento foi reivindicado. --- .../out/messaging/OutboxClaimCoordinator.java | 53 +++++++ .../adapter/out/messaging/OutboxRelay.java | 44 +----- .../persistence/outbox/OutboxEventEntity.java | 9 ++ .../outbox/OutboxEventJpaRepository.java | 14 ++ .../V4__outbox_events_add_claimed_at.sql | 1 + .../messaging/OutboxClaimCoordinatorTest.java | 137 ++++++++++++++++++ ...OutboxClaimCoordinatorTransactionalIT.java | 65 +++++++++ .../out/messaging/OutboxRelayTest.java | 59 ++------ 8 files changed, 301 insertions(+), 81 deletions(-) create mode 100644 src/main/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinator.java create mode 100644 src/main/resources/db/migration/V4__outbox_events_add_claimed_at.sql create mode 100644 src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTest.java create mode 100644 src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTransactionalIT.java diff --git a/src/main/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinator.java b/src/main/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinator.java new file mode 100644 index 0000000..480e2b8 --- /dev/null +++ b/src/main/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinator.java @@ -0,0 +1,53 @@ +package com.lmoraesdev.payment.adapter.out.messaging; + +import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventEntity; +import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventJpaRepository; +import java.time.Duration; +import java.time.Instant; +import java.util.List; +import java.util.UUID; +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Transactional; + +@Component +public class OutboxClaimCoordinator { + + private static final Duration STUCK_IN_FLIGHT_THRESHOLD = Duration.ofMinutes(2); + + private final OutboxEventJpaRepository repository; + + public OutboxClaimCoordinator(OutboxEventJpaRepository repository) { + this.repository = repository; + } + + @Transactional + public List claimBatch() { + repository.reapStuckInFlight(Instant.now().minus(STUCK_IN_FLIGHT_THRESHOLD)); + + List claimed = repository.findBatchForUpdateSkipLocked(); + claimed.forEach(OutboxEventEntity::markInFlight); + return repository.saveAll(claimed); + } + + @Transactional + public void markPublished(UUID eventId) { + repository + .findById(eventId) + .ifPresent( + event -> { + event.markPublished(); + repository.save(event); + }); + } + + @Transactional + public void revertToPending(UUID eventId) { + repository + .findById(eventId) + .ifPresent( + event -> { + event.revertToPending(); + repository.save(event); + }); + } +} diff --git a/src/main/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxRelay.java b/src/main/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxRelay.java index 1e07e06..c771eea 100644 --- a/src/main/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxRelay.java +++ b/src/main/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxRelay.java @@ -1,7 +1,6 @@ package com.lmoraesdev.payment.adapter.out.messaging; import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventEntity; -import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventJpaRepository; import com.lmoraesdev.payment.config.logging.Logger5w1hBuilder; import io.micrometer.core.instrument.Counter; import io.micrometer.core.instrument.MeterRegistry; @@ -9,29 +8,27 @@ import java.time.Duration; import java.time.Instant; import java.util.List; -import java.util.UUID; import org.slf4j.MDC; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; -import org.springframework.transaction.annotation.Transactional; @Component public class OutboxRelay { private static final String TOPIC = "payments.charge-created"; - private final OutboxEventJpaRepository repository; + private final OutboxClaimCoordinator outboxClaimCoordinator; private final KafkaTemplate kafkaTemplate; private final Counter publishedCounter; private final Counter failedCounter; private final Timer publishLagTimer; public OutboxRelay( - OutboxEventJpaRepository repository, + OutboxClaimCoordinator outboxClaimCoordinator, KafkaTemplate kafkaTemplate, MeterRegistry meterRegistry) { - this.repository = repository; + this.outboxClaimCoordinator = outboxClaimCoordinator; this.kafkaTemplate = kafkaTemplate; this.publishedCounter = meterRegistry.counter("outbox_events_published_total"); this.failedCounter = meterRegistry.counter("outbox_events_failed_total"); @@ -40,30 +37,23 @@ public OutboxRelay( @Scheduled(fixedDelay = 5000) public void publishPending() { - List claimed = claimBatch(); + List claimed = outboxClaimCoordinator.claimBatch(); for (OutboxEventEntity event : claimed) { publish(event); } } - @Transactional - public List claimBatch() { - List claimed = repository.findBatchForUpdateSkipLocked(); - claimed.forEach(OutboxEventEntity::markInFlight); - return repository.saveAll(claimed); - } - private void publish(OutboxEventEntity event) { try { kafkaTemplate.send(TOPIC, event.getAggregateId(), event.getPayload()).get(); - markPublished(event.getId()); + outboxClaimCoordinator.markPublished(event.getId()); publishedCounter.increment(); publishLagTimer.record(Duration.between(event.getCreatedAt(), Instant.now())); } catch (Exception e) { failedCounter.increment(); logPublishFailure(event, e); - revertToPending(event.getId()); + outboxClaimCoordinator.revertToPending(event.getId()); } } @@ -86,26 +76,4 @@ private void logPublishFailure(OutboxEventEntity event, Exception e) { } } } - - @Transactional - public void markPublished(UUID eventId) { - repository - .findById(eventId) - .ifPresent( - event -> { - event.markPublished(); - repository.save(event); - }); - } - - @Transactional - public void revertToPending(UUID eventId) { - repository - .findById(eventId) - .ifPresent( - event -> { - event.revertToPending(); - repository.save(event); - }); - } } diff --git a/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/outbox/OutboxEventEntity.java b/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/outbox/OutboxEventEntity.java index 73b4089..c1670b7 100644 --- a/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/outbox/OutboxEventEntity.java +++ b/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/outbox/OutboxEventEntity.java @@ -46,6 +46,9 @@ public class OutboxEventEntity { @Column(name = "correlation_id") private String correlationId; + @Column(name = "claimed_at") + private Instant claimedAt; + protected OutboxEventEntity() { // JPA } @@ -70,6 +73,7 @@ public static OutboxEventEntity pending( public void markInFlight() { this.status = OutboxStatus.IN_FLIGHT; + this.claimedAt = Instant.now(); } public void markPublished() { @@ -79,6 +83,7 @@ public void markPublished() { public void revertToPending() { this.status = OutboxStatus.PENDING; + this.claimedAt = null; } public UUID getId() { @@ -108,4 +113,8 @@ public Instant getCreatedAt() { public String getCorrelationId() { return correlationId; } + + public Instant getClaimedAt() { + return claimedAt; + } } diff --git a/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/outbox/OutboxEventJpaRepository.java b/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/outbox/OutboxEventJpaRepository.java index 096d8ca..20cdf7e 100644 --- a/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/outbox/OutboxEventJpaRepository.java +++ b/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/outbox/OutboxEventJpaRepository.java @@ -1,9 +1,12 @@ package com.lmoraesdev.payment.adapter.out.persistence.outbox; +import java.time.Instant; import java.util.List; import java.util.UUID; import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Modifying; import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; public interface OutboxEventJpaRepository extends JpaRepository { @@ -18,4 +21,15 @@ public interface OutboxEventJpaRepository extends JpaRepository findBatchForUpdateSkipLocked(); + + @Modifying + @Query( + value = + """ + update outbox_events + set status = 'PENDING', claimed_at = null + where status = 'IN_FLIGHT' and claimed_at < :cutoff + """, + nativeQuery = true) + int reapStuckInFlight(@Param("cutoff") Instant cutoff); } diff --git a/src/main/resources/db/migration/V4__outbox_events_add_claimed_at.sql b/src/main/resources/db/migration/V4__outbox_events_add_claimed_at.sql new file mode 100644 index 0000000..e56e78c --- /dev/null +++ b/src/main/resources/db/migration/V4__outbox_events_add_claimed_at.sql @@ -0,0 +1 @@ +alter table outbox_events add column claimed_at timestamptz; diff --git a/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTest.java b/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTest.java new file mode 100644 index 0000000..aa05b9d --- /dev/null +++ b/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTest.java @@ -0,0 +1,137 @@ +package com.lmoraesdev.payment.adapter.out.messaging; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventEntity; +import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventJpaRepository; +import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxStatus; +import java.time.Instant; +import java.util.List; +import java.util.Optional; +import java.util.UUID; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Nested; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InOrder; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.junit.jupiter.MockitoExtension; + +@DisplayName("OutboxClaimCoordinator") +@ExtendWith(MockitoExtension.class) +class OutboxClaimCoordinatorTest { + + @Mock OutboxEventJpaRepository repository; + + OutboxClaimCoordinator coordinator; + + @Nested + @DisplayName("claimBatch") + class ClaimBatch { + + @Test + @DisplayName("reclama eventos IN_FLIGHT travados antes de buscar o próximo lote PENDING") + void reapsStuckInFlightBeforeFetchingPendingBatch() { + coordinator = new OutboxClaimCoordinator(repository); + when(repository.reapStuckInFlight(any(Instant.class))).thenReturn(0); + when(repository.findBatchForUpdateSkipLocked()).thenReturn(List.of()); + when(repository.saveAll(List.of())).thenReturn(List.of()); + + coordinator.claimBatch(); + + InOrder inOrder = Mockito.inOrder(repository); + inOrder.verify(repository).reapStuckInFlight(any(Instant.class)); + inOrder.verify(repository).findBatchForUpdateSkipLocked(); + } + + @Test + @DisplayName("marca o lote reivindicado como IN_FLIGHT com claimedAt e persiste") + void marksFetchedBatchInFlightAndPersists() { + coordinator = new OutboxClaimCoordinator(repository); + OutboxEventEntity event = + OutboxEventEntity.pending("Charge", "charge-1", "ChargeCreated", "{}", null); + when(repository.reapStuckInFlight(any(Instant.class))).thenReturn(0); + when(repository.findBatchForUpdateSkipLocked()).thenReturn(List.of(event)); + when(repository.saveAll(List.of(event))) + .thenAnswer( + invocation -> { + assertThat(event.getStatus()).isEqualTo(OutboxStatus.IN_FLIGHT); + assertThat(event.getClaimedAt()).isNotNull(); + return List.of(event); + }); + + List claimed = coordinator.claimBatch(); + + assertThat(claimed).containsExactly(event); + } + } + + @Nested + @DisplayName("markPublished") + class MarkPublished { + + @Test + @DisplayName("marca o evento existente como PUBLISHED e persiste") + void marksExistingEventPublished() { + coordinator = new OutboxClaimCoordinator(repository); + OutboxEventEntity event = + OutboxEventEntity.pending("Charge", "charge-2", "ChargeCreated", "{}", null); + when(repository.findById(event.getId())).thenReturn(Optional.of(event)); + + coordinator.markPublished(event.getId()); + + assertThat(event.getStatus()).isEqualTo(OutboxStatus.PUBLISHED); + verify(repository).save(event); + } + + @Test + @DisplayName("não faz nada quando o evento não é encontrado") + void doesNothingWhenEventNotFound() { + coordinator = new OutboxClaimCoordinator(repository); + UUID missingId = UUID.randomUUID(); + when(repository.findById(missingId)).thenReturn(Optional.empty()); + + coordinator.markPublished(missingId); + + verify(repository, never()).save(any()); + } + } + + @Nested + @DisplayName("revertToPending") + class RevertToPending { + + @Test + @DisplayName("reverte o evento existente para PENDING e limpa claimedAt") + void revertsExistingEventToPendingAndClearsClaimedAt() { + coordinator = new OutboxClaimCoordinator(repository); + OutboxEventEntity event = + OutboxEventEntity.pending("Charge", "charge-3", "ChargeCreated", "{}", null); + event.markInFlight(); + when(repository.findById(event.getId())).thenReturn(Optional.of(event)); + + coordinator.revertToPending(event.getId()); + + assertThat(event.getStatus()).isEqualTo(OutboxStatus.PENDING); + assertThat(event.getClaimedAt()).isNull(); + verify(repository).save(event); + } + + @Test + @DisplayName("não faz nada quando o evento não é encontrado") + void doesNothingWhenEventNotFound() { + coordinator = new OutboxClaimCoordinator(repository); + UUID missingId = UUID.randomUUID(); + when(repository.findById(missingId)).thenReturn(Optional.empty()); + + coordinator.revertToPending(missingId); + + verify(repository, never()).save(any()); + } + } +} diff --git a/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTransactionalIT.java b/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTransactionalIT.java new file mode 100644 index 0000000..b07ce96 --- /dev/null +++ b/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTransactionalIT.java @@ -0,0 +1,65 @@ +package com.lmoraesdev.payment.adapter.out.messaging; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.anyList; +import static org.mockito.Mockito.doAnswer; + +import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventEntity; +import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventJpaRepository; +import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxStatus; +import com.lmoraesdev.payment.support.AbstractIntegrationTest; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Import; +import org.springframework.test.context.bean.override.mockito.MockitoSpyBean; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +// Propagation.NOT_SUPPORTED desliga o wrapper transacional do próprio teste, senão o rollback +// do @Transactional real de claimBatch() ficaria preso na mesma transação nunca-commitada do teste. +@DisplayName("OutboxClaimCoordinator (transacional real)") +@Import(OutboxClaimCoordinator.class) +@Transactional(propagation = Propagation.NOT_SUPPORTED) +class OutboxClaimCoordinatorTransactionalIT extends AbstractIntegrationTest { + + @Autowired OutboxClaimCoordinator coordinator; + + @MockitoSpyBean OutboxEventJpaRepository repository; + + private OutboxEventEntity seeded; + + @AfterEach + void cleanUp() { + if (seeded != null) { + repository.deleteById(seeded.getId()); + } + } + + @Test + @DisplayName("exceção no meio de claimBatch desfaz a mudança de status pro evento reivindicado") + void exceptionMidClaimBatchRollsBackStatusChange() { + seeded = + repository.save( + OutboxEventEntity.pending( + "Charge", "aggregate-rollback", "ChargeCreated", "{}", null)); + + doAnswer( + invocation -> { + invocation.callRealMethod(); + throw new RuntimeException("boom-mid-claim"); + }) + .when(repository) + .saveAll(anyList()); + + assertThatThrownBy(() -> coordinator.claimBatch()) + .isInstanceOf(RuntimeException.class) + .hasMessage("boom-mid-claim"); + + OutboxEventEntity reloaded = repository.findById(seeded.getId()).orElseThrow(); + assertThat(reloaded.getStatus()).isEqualTo(OutboxStatus.PENDING); + assertThat(reloaded.getClaimedAt()).isNull(); + } +} diff --git a/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxRelayTest.java b/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxRelayTest.java index 6e3e2ea..8a5ecc0 100644 --- a/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxRelayTest.java +++ b/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxRelayTest.java @@ -10,11 +10,8 @@ import ch.qos.logback.classic.spi.ILoggingEvent; import ch.qos.logback.core.read.ListAppender; import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventEntity; -import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventJpaRepository; -import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxStatus; import io.micrometer.core.instrument.simple.SimpleMeterRegistry; import java.util.List; -import java.util.Optional; import java.util.concurrent.CompletableFuture; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; @@ -32,7 +29,7 @@ @ExtendWith(MockitoExtension.class) class OutboxRelayTest { - @Mock OutboxEventJpaRepository repository; + @Mock OutboxClaimCoordinator outboxClaimCoordinator; @Mock KafkaTemplate kafkaTemplate; @@ -46,7 +43,7 @@ class OutboxRelayTest { @BeforeEach void setUp() { meterRegistry = new SimpleMeterRegistry(); - relay = new OutboxRelay(repository, kafkaTemplate, meterRegistry); + relay = new OutboxRelay(outboxClaimCoordinator, kafkaTemplate, meterRegistry); logger = (Logger) LoggerFactory.getLogger(OutboxRelay.class); appender.start(); logger.addAppender(appender); @@ -58,60 +55,38 @@ void tearDown() { } @Test - @DisplayName("reivindica lote PENDING como IN_FLIGHT, publica e marca PUBLISHED") + @DisplayName("reivindica lote via OutboxClaimCoordinator, publica e marca PUBLISHED") void claimsPublishesAndMarksPublished() { OutboxEventEntity event = OutboxEventEntity.pending("Charge", "charge-1", "ChargeCreated", "{}", null); - when(repository.findBatchForUpdateSkipLocked()).thenReturn(List.of(event)); - when(repository.saveAll(List.of(event))).thenReturn(List.of(event)); - when(repository.findById(event.getId())).thenReturn(Optional.of(event)); + when(outboxClaimCoordinator.claimBatch()).thenReturn(List.of(event)); when(kafkaTemplate.send(any(String.class), any(), any())) .thenReturn(CompletableFuture.completedFuture(mockSendResult())); relay.publishPending(); - assertThat(event.getStatus()).isEqualTo(OutboxStatus.PUBLISHED); - verify(repository).saveAll(List.of(event)); - verify(repository).save(event); + verify(outboxClaimCoordinator).markPublished(event.getId()); + verify(outboxClaimCoordinator, never()).revertToPending(any()); assertThat(meterRegistry.counter("outbox_events_published_total").count()).isEqualTo(1.0); assertThat(meterRegistry.counter("outbox_events_failed_total").count()).isEqualTo(0.0); assertThat(meterRegistry.timer("outbox_publish_lag").count()).isEqualTo(1L); } @Test - @DisplayName("claimBatch marca o lote reivindicado como IN_FLIGHT antes de publicar") - void claimBatchMarksEventsInFlight() { - OutboxEventEntity event = - OutboxEventEntity.pending("Charge", "charge-2", "ChargeCreated", "{}", null); - when(repository.findBatchForUpdateSkipLocked()).thenReturn(List.of(event)); - when(repository.saveAll(List.of(event))) - .thenAnswer( - inv -> { - assertThat(event.getStatus()).isEqualTo(OutboxStatus.IN_FLIGHT); - return List.of(event); - }); - - List claimed = relay.claimBatch(); - - assertThat(claimed).containsExactly(event); - } - - @Test - @DisplayName("falha ao publicar loga via Logger5w1hBuilder e reverte pra PENDING") + @DisplayName( + "falha ao publicar loga via Logger5w1hBuilder e reverte pra PENDING via coordinator") void logsFailureAndRevertsToPending() { OutboxEventEntity event = OutboxEventEntity.pending( "Charge", "charge-3", "ChargeCreated", "{}", "trace-original-request"); - when(repository.findBatchForUpdateSkipLocked()).thenReturn(List.of(event)); - when(repository.saveAll(List.of(event))).thenReturn(List.of(event)); - when(repository.findById(event.getId())).thenReturn(Optional.of(event)); + when(outboxClaimCoordinator.claimBatch()).thenReturn(List.of(event)); when(kafkaTemplate.send(any(String.class), any(), any())) .thenReturn(CompletableFuture.failedFuture(new RuntimeException("kafka down"))); relay.publishPending(); - assertThat(event.getStatus()).isEqualTo(OutboxStatus.PENDING); - verify(repository).save(event); + verify(outboxClaimCoordinator).revertToPending(event.getId()); + verify(outboxClaimCoordinator, never()).markPublished(any()); assertThat(appender.list).hasSize(1); ILoggingEvent logged = appender.list.get(0); @@ -129,9 +104,7 @@ void logsFailureAndRevertsToPending() { void doesNotTouchMdcWhenCorrelationIdIsAbsent() { OutboxEventEntity event = OutboxEventEntity.pending("Charge", "charge-6", "ChargeCreated", "{}", null); - when(repository.findBatchForUpdateSkipLocked()).thenReturn(List.of(event)); - when(repository.saveAll(List.of(event))).thenReturn(List.of(event)); - when(repository.findById(event.getId())).thenReturn(Optional.of(event)); + when(outboxClaimCoordinator.claimBatch()).thenReturn(List.of(event)); when(kafkaTemplate.send(any(String.class), any(), any())) .thenReturn(CompletableFuture.failedFuture(new RuntimeException("kafka down"))); @@ -143,15 +116,15 @@ void doesNotTouchMdcWhenCorrelationIdIsAbsent() { } @Test - @DisplayName("lote vazio não publica nem salva nada") + @DisplayName("lote vazio não publica nem marca nada") void doesNothingWhenBatchIsEmpty() { - when(repository.findBatchForUpdateSkipLocked()).thenReturn(List.of()); - when(repository.saveAll(List.of())).thenReturn(List.of()); + when(outboxClaimCoordinator.claimBatch()).thenReturn(List.of()); relay.publishPending(); verify(kafkaTemplate, never()).send(any(String.class), any(), any()); - verify(repository, never()).save(any()); + verify(outboxClaimCoordinator, never()).markPublished(any()); + verify(outboxClaimCoordinator, never()).revertToPending(any()); } @SuppressWarnings("unchecked") From aa30426aa8ac2ab2b00e4008aa5741ccbe12d9e5 Mon Sep 17 00:00:00 2001 From: Leandro Moraes Date: Thu, 23 Jul 2026 10:11:09 -0300 Subject: [PATCH 2/4] =?UTF-8?q?fix(charges):=20otimistic=20locking,=20vali?= =?UTF-8?q?da=C3=A7=C3=A3o=20defensiva=20no=20webhook,=20e=20replay=20corr?= =?UTF-8?q?eto=20sob=20concorr=C3=AAncia=20na=20idempot=C3=AAncia?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adiciona @Version (coluna version, migration V5) em charges. ProcessWebhookService e ChargeExpirationJob podiam sobrescrever a mesma charge concorrentemente sem avisar ninguém. No webhook, o conflito agora propaga OptimisticLockingFailureException até o GlobalExceptionHandler, que responde 409 (o provedor reenvia; dedup por event_id garante que reprocessar é seguro). No job de expiração, extrai ChargeExpirationCoordinator (@Transactional por charge, evitando o self-invocation que anularia a anotação) para que uma charge em conflito seja pulada nesse ciclo sem impedir as demais de serem processadas. ProcessWebhookService agora valida o status recebido antes de qualquer transição, lançando InvalidChargeStatusException (422) em vez de deixar um IllegalArgumentException cru virar 500. CreateChargeService: duas requisições concorrentes com a mesma Idempotency-Key nova faziam a perdedora receber um 500 da violação de unicidade, mesmo com o rollback já evitando duplicar dados. Extrai ChargeCreationCoordinator (@Transactional) para a escrita; ao capturar DataIntegrityViolationException UMA CAMADA ACIMA da transação (depois do rollback completo do Postgres), CreateChargeService busca o registro de idempotência da vencedora numa transação nova e retorna o replay em vez do erro. Testes cobrem os três cenários; a concorrência real (duas threads com a mesma key) é provada via Testcontainers em CreateChargeServiceConcurrencyIT. --- .../in/web/GlobalExceptionHandler.java | 21 ++++ .../out/persistence/ChargeJpaEntity.java | 13 +++ .../adapter/out/persistence/ChargeMapper.java | 4 +- .../ChargeExpirationCoordinator.java | 51 +++++++++ .../out/scheduling/ChargeExpirationJob.java | 47 +++----- .../usecase/ChargeCreationCoordinator.java | 71 ++++++++++++ .../usecase/CreateChargeService.java | 73 ++++-------- .../usecase/ProcessWebhookService.java | 11 +- .../InvalidChargeStatusException.java | 8 ++ .../payment/domain/model/Charge.java | 25 +++- .../db/migration/V5__charges_add_version.sql | 1 + .../in/web/GlobalExceptionHandlerTest.java | 24 ++++ .../out/persistence/ChargeRepositoryIT.java | 4 +- .../ChargeExpirationCoordinatorTest.java | 48 ++++++++ .../scheduling/ChargeExpirationJobTest.java | 40 ++++--- .../ChargeCreationCoordinatorTest.java | 86 ++++++++++++++ .../CreateChargeServiceConcurrencyIT.java | 71 ++++++++++++ .../usecase/CreateChargeServiceTest.java | 107 ++++++++++-------- .../usecase/ProcessWebhookServiceTest.java | 41 +++++++ .../payment/domain/model/ChargeTest.java | 13 ++- .../payment/testdata/ChargeTestData.java | 8 +- 21 files changed, 615 insertions(+), 152 deletions(-) create mode 100644 src/main/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationCoordinator.java create mode 100644 src/main/java/com/lmoraesdev/payment/application/usecase/ChargeCreationCoordinator.java create mode 100644 src/main/java/com/lmoraesdev/payment/domain/exception/InvalidChargeStatusException.java create mode 100644 src/main/resources/db/migration/V5__charges_add_version.sql create mode 100644 src/test/java/com/lmoraesdev/payment/adapter/in/web/GlobalExceptionHandlerTest.java create mode 100644 src/test/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationCoordinatorTest.java create mode 100644 src/test/java/com/lmoraesdev/payment/application/usecase/ChargeCreationCoordinatorTest.java create mode 100644 src/test/java/com/lmoraesdev/payment/application/usecase/CreateChargeServiceConcurrencyIT.java diff --git a/src/main/java/com/lmoraesdev/payment/adapter/in/web/GlobalExceptionHandler.java b/src/main/java/com/lmoraesdev/payment/adapter/in/web/GlobalExceptionHandler.java index a557629..d7a931f 100644 --- a/src/main/java/com/lmoraesdev/payment/adapter/in/web/GlobalExceptionHandler.java +++ b/src/main/java/com/lmoraesdev/payment/adapter/in/web/GlobalExceptionHandler.java @@ -8,6 +8,7 @@ import org.slf4j.MDC; import org.springframework.core.Ordered; import org.springframework.core.annotation.Order; +import org.springframework.dao.OptimisticLockingFailureException; import org.springframework.http.HttpStatus; import org.springframework.http.ProblemDetail; import org.springframework.web.bind.MethodArgumentNotValidException; @@ -67,6 +68,26 @@ public ProblemDetail handleDomain(DomainException ex) { return problem; } + // 409 — conflito de concorrência otimista. Recuperável: cliente/provedor pode tentar de novo. + @ExceptionHandler(OptimisticLockingFailureException.class) + public ProblemDetail handleOptimisticLock(OptimisticLockingFailureException ex) { + Logger5w1hBuilder.create(GlobalExceptionHandler.class) + .where("GlobalExceptionHandler") + .what("optimistic_lock_conflict") + .why("recurso modificado concorrentemente, quem chamou deve tentar de novo") + .who("system") + .how("exception handling") + .error(ex); + + ProblemDetail problem = + ProblemDetail.forStatusAndDetail( + HttpStatus.CONFLICT, + "Recurso modificado concorrentemente, tente novamente"); + problem.setTitle("Concurrent modification conflict"); + addTraceId(problem); + return problem; + } + // 500 — inesperado. AQUI sim loga, em ERROR, com a stack. @ExceptionHandler(Exception.class) public ProblemDetail handleUnexpected(Exception ex) { diff --git a/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeJpaEntity.java b/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeJpaEntity.java index c86fae1..59d3b4e 100644 --- a/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeJpaEntity.java +++ b/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeJpaEntity.java @@ -7,6 +7,7 @@ import jakarta.persistence.Enumerated; import jakarta.persistence.Id; import jakarta.persistence.Table; +import jakarta.persistence.Version; import java.time.Instant; import java.util.UUID; @@ -28,6 +29,10 @@ public class ChargeJpaEntity { @Column(name = "expires_at", nullable = false) private Instant expiresAt; + @Version + @Column(nullable = false) + private Long version; + protected ChargeJpaEntity() {} public UUID getId() { @@ -69,4 +74,12 @@ public Instant getExpiresAt() { public void setExpiresAt(Instant expiresAt) { this.expiresAt = expiresAt; } + + public Long getVersion() { + return version; + } + + public void setVersion(Long version) { + this.version = version; + } } diff --git a/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeMapper.java b/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeMapper.java index 4a0a0ef..5f770f0 100644 --- a/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeMapper.java +++ b/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeMapper.java @@ -16,6 +16,7 @@ static ChargeJpaEntity toEntity(Charge charge) { entity.setStatus(charge.getStatus()); entity.setCreatedAt(charge.getCreatedAt()); entity.setExpiresAt(charge.getExpiresAt()); + entity.setVersion(charge.getVersion()); return entity; } @@ -26,6 +27,7 @@ static Charge toDomain(ChargeJpaEntity entity) { new Money(BigDecimal.valueOf(entity.getAmountCentavos(), 2)), entity.getStatus(), entity.getCreatedAt(), - entity.getExpiresAt()); + entity.getExpiresAt(), + entity.getVersion()); } } diff --git a/src/main/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationCoordinator.java b/src/main/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationCoordinator.java new file mode 100644 index 0000000..7a1d3ed --- /dev/null +++ b/src/main/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationCoordinator.java @@ -0,0 +1,51 @@ +package com.lmoraesdev.payment.adapter.out.scheduling; + +import com.lmoraesdev.payment.application.port.out.ChargeRepository; +import com.lmoraesdev.payment.application.port.out.OutboxEventPort; +import com.lmoraesdev.payment.config.logging.Logger5w1hBuilder; +import com.lmoraesdev.payment.domain.event.ChargeStatusChangedEvent; +import com.lmoraesdev.payment.domain.model.Charge; +import com.lmoraesdev.payment.domain.model.ChargeStatus; +import java.time.Instant; +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Transactional; + +@Component +public class ChargeExpirationCoordinator { + + private final ChargeRepository chargeRepository; + private final OutboxEventPort outboxEventPort; + + public ChargeExpirationCoordinator( + ChargeRepository chargeRepository, OutboxEventPort outboxEventPort) { + this.chargeRepository = chargeRepository; + this.outboxEventPort = outboxEventPort; + } + + @Transactional + public void expireOne(Charge charge) { + ChargeStatus previousStatus = charge.getStatus(); + + charge.transitionTo(ChargeStatus.EXPIRED); + + chargeRepository.save(charge); + + outboxEventPort.record( + "Charge", + charge.getId().toString(), + "ChargeExpired", + new ChargeStatusChangedEvent( + charge.getId(), + previousStatus.name(), + ChargeStatus.EXPIRED.name(), + Instant.now())); + + Logger5w1hBuilder.create(ChargeExpirationCoordinator.class) + .where("ChargeExpirationCoordinator") + .what("charge_state_transitioned") + .why("ttl_expired") + .who("system") + .how("scheduled_expiration") + .info(); + } +} diff --git a/src/main/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationJob.java b/src/main/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationJob.java index 1cc333b..1919bc5 100644 --- a/src/main/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationJob.java +++ b/src/main/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationJob.java @@ -1,61 +1,50 @@ package com.lmoraesdev.payment.adapter.out.scheduling; import com.lmoraesdev.payment.application.port.out.ChargeRepository; -import com.lmoraesdev.payment.application.port.out.OutboxEventPort; import com.lmoraesdev.payment.config.logging.Logger5w1hBuilder; -import com.lmoraesdev.payment.domain.event.ChargeStatusChangedEvent; import com.lmoraesdev.payment.domain.model.Charge; -import com.lmoraesdev.payment.domain.model.ChargeStatus; import java.time.Instant; import java.util.List; +import org.springframework.dao.OptimisticLockingFailureException; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; -import org.springframework.transaction.annotation.Transactional; @Component public class ChargeExpirationJob { private final ChargeRepository chargeRepository; - private final OutboxEventPort outboxEventPort; + private final ChargeExpirationCoordinator chargeExpirationCoordinator; - public ChargeExpirationJob(ChargeRepository chargeRepository, OutboxEventPort outboxEventPort) { + public ChargeExpirationJob( + ChargeRepository chargeRepository, + ChargeExpirationCoordinator chargeExpirationCoordinator) { this.chargeRepository = chargeRepository; - this.outboxEventPort = outboxEventPort; + this.chargeExpirationCoordinator = chargeExpirationCoordinator; } @Scheduled(fixedDelay = 60000) - @Transactional public void expireOverdueCharges() { List expired = chargeRepository.findExpiredActive(Instant.now()); for (Charge charge : expired) { - expire(charge); + try { + chargeExpirationCoordinator.expireOne(charge); + } catch (OptimisticLockingFailureException e) { + logSkippedDueToConflict(charge, e); + } } } - private void expire(Charge charge) { - ChargeStatus previousStatus = charge.getStatus(); - - charge.transitionTo(ChargeStatus.EXPIRED); - - chargeRepository.save(charge); - - outboxEventPort.record( - "Charge", - charge.getId().toString(), - "ChargeExpired", - new ChargeStatusChangedEvent( - charge.getId(), - previousStatus.name(), - ChargeStatus.EXPIRED.name(), - Instant.now())); - + private void logSkippedDueToConflict(Charge charge, OptimisticLockingFailureException e) { Logger5w1hBuilder.create(ChargeExpirationJob.class) .where("ChargeExpirationJob") - .what("charge_state_transitioned") - .why("ttl_expired") + .what("charge_expiration_skipped") + .why( + "optimistic lock conflict on charge " + + charge.getId() + + ", updated concorrentemente, tentará de novo no próximo poll") .who("system") .how("scheduled_expiration") - .info(); + .error(e); } } diff --git a/src/main/java/com/lmoraesdev/payment/application/usecase/ChargeCreationCoordinator.java b/src/main/java/com/lmoraesdev/payment/application/usecase/ChargeCreationCoordinator.java new file mode 100644 index 0000000..2e531b4 --- /dev/null +++ b/src/main/java/com/lmoraesdev/payment/application/usecase/ChargeCreationCoordinator.java @@ -0,0 +1,71 @@ +package com.lmoraesdev.payment.application.usecase; + +import com.lmoraesdev.payment.application.port.in.CreateChargeCommand; +import com.lmoraesdev.payment.application.port.in.CreateChargeResult; +import com.lmoraesdev.payment.application.port.out.ChargeRepository; +import com.lmoraesdev.payment.application.port.out.IdempotencyPort; +import com.lmoraesdev.payment.application.port.out.OutboxEventPort; +import com.lmoraesdev.payment.config.logging.Logger5w1hBuilder; +import com.lmoraesdev.payment.domain.event.ChargeCreatedEvent; +import com.lmoraesdev.payment.domain.model.Charge; +import com.lmoraesdev.payment.domain.model.Money; +import io.micrometer.core.instrument.Counter; +import io.micrometer.core.instrument.MeterRegistry; +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Transactional; + +@Component +public class ChargeCreationCoordinator { + + private final ChargeRepository chargeRepository; + private final OutboxEventPort outboxEventPort; + private final IdempotencyPort idempotencyPort; + private final Counter chargesCreatedCounter; + + public ChargeCreationCoordinator( + ChargeRepository chargeRepository, + OutboxEventPort outboxEventPort, + IdempotencyPort idempotencyPort, + MeterRegistry meterRegistry) { + this.chargeRepository = chargeRepository; + this.outboxEventPort = outboxEventPort; + this.idempotencyPort = idempotencyPort; + this.chargesCreatedCounter = meterRegistry.counter("charges_created_total"); + } + + @Transactional + public CreateChargeResult createAndPersist( + CreateChargeCommand command, Money amount, String requestHash) { + Charge charge = Charge.create(amount); + + Charge saved = chargeRepository.save(charge); + + outboxEventPort.record( + "Charge", + saved.getId().toString(), + "ChargeCreated", + new ChargeCreatedEvent( + saved.getId(), saved.getAmount().amount(), saved.getCreatedAt())); + + Logger5w1hBuilder.create(ChargeCreationCoordinator.class) + .where("ChargeCreationCoordinator") + .what("charge_created") + .why("charge creation requested") + .who("system") + .how("createCharge") + .info(); + + CreateChargeResult result = + new CreateChargeResult( + saved.getId(), + saved.getStatus().name(), + saved.getAmount().amount(), + saved.getCreatedAt(), + false); + + idempotencyPort.save(command.idempotencyKey(), requestHash, saved.getId(), result); + chargesCreatedCounter.increment(); + + return result; + } +} diff --git a/src/main/java/com/lmoraesdev/payment/application/usecase/CreateChargeService.java b/src/main/java/com/lmoraesdev/payment/application/usecase/CreateChargeService.java index 43b4134..0100b2b 100644 --- a/src/main/java/com/lmoraesdev/payment/application/usecase/CreateChargeService.java +++ b/src/main/java/com/lmoraesdev/payment/application/usecase/CreateChargeService.java @@ -3,48 +3,32 @@ import com.lmoraesdev.payment.application.port.in.CreateCharge; import com.lmoraesdev.payment.application.port.in.CreateChargeCommand; import com.lmoraesdev.payment.application.port.in.CreateChargeResult; -import com.lmoraesdev.payment.application.port.out.ChargeRepository; import com.lmoraesdev.payment.application.port.out.IdempotencyPort; import com.lmoraesdev.payment.application.port.out.IdempotencyPort.StoredIdempotency; -import com.lmoraesdev.payment.application.port.out.OutboxEventPort; -import com.lmoraesdev.payment.config.logging.Logger5w1hBuilder; -import com.lmoraesdev.payment.domain.event.ChargeCreatedEvent; import com.lmoraesdev.payment.domain.exception.IdempotencyConflictException; -import com.lmoraesdev.payment.domain.model.Charge; import com.lmoraesdev.payment.domain.model.Money; -import io.micrometer.core.instrument.Counter; -import io.micrometer.core.instrument.MeterRegistry; import java.math.BigDecimal; import java.nio.charset.StandardCharsets; import java.security.MessageDigest; import java.security.NoSuchAlgorithmException; import java.util.HexFormat; import java.util.Optional; +import org.springframework.dao.DataIntegrityViolationException; import org.springframework.stereotype.Service; -import org.springframework.transaction.annotation.Transactional; @Service public class CreateChargeService implements CreateCharge { - private final ChargeRepository chargeRepository; - private final OutboxEventPort outboxEventPort; + private final ChargeCreationCoordinator chargeCreationCoordinator; private final IdempotencyPort idempotencyPort; - private final Counter chargesCreatedCounter; public CreateChargeService( - ChargeRepository chargeRepository, - OutboxEventPort outboxEventPort, - IdempotencyPort idempotencyPort, - MeterRegistry meterRegistry) { - this.chargeRepository = chargeRepository; - this.outboxEventPort = outboxEventPort; + ChargeCreationCoordinator chargeCreationCoordinator, IdempotencyPort idempotencyPort) { + this.chargeCreationCoordinator = chargeCreationCoordinator; this.idempotencyPort = idempotencyPort; - this.chargesCreatedCounter = meterRegistry.counter("charges_created_total"); } @Override - @Transactional public CreateChargeResult create(CreateChargeCommand command) { - Money amount = new Money(command.amount()); String requestHash = hash(command.amount()); @@ -53,37 +37,24 @@ public CreateChargeResult create(CreateChargeCommand command) { return replay(existing.get(), requestHash, command.idempotencyKey()); } - Charge charge = Charge.create(amount); - - Charge saved = chargeRepository.save(charge); - - outboxEventPort.record( - "Charge", - saved.getId().toString(), - "ChargeCreated", - new ChargeCreatedEvent( - saved.getId(), saved.getAmount().amount(), saved.getCreatedAt())); - - Logger5w1hBuilder.create(CreateChargeService.class) - .where("CreateChargeService") - .what("charge_created") - .why("charge creation requested") - .who("system") - .how("createCharge") - .info(); - - CreateChargeResult result = - new CreateChargeResult( - saved.getId(), - saved.getStatus().name(), - saved.getAmount().amount(), - saved.getCreatedAt(), - false); - - idempotencyPort.save(command.idempotencyKey(), requestHash, saved.getId(), result); - chargesCreatedCounter.increment(); - - return result; + try { + return chargeCreationCoordinator.createAndPersist(command, amount, requestHash); + } catch (DataIntegrityViolationException e) { + // A transação da tentativa de criação já sofreu rollback completo (violação de + // constraint aborta a transação inteira no Postgres). A requisição vencedora já + // deve ter commitado seu registro de idempotência; buscamos numa transação nova. + StoredIdempotency winner = + idempotencyPort + .findByKey(command.idempotencyKey()) + .orElseThrow( + () -> + new IllegalStateException( + "registro de idempotência esperado após" + + " conflito de constraint não" + + " encontrado para key: " + + command.idempotencyKey())); + return replay(winner, requestHash, command.idempotencyKey()); + } } private CreateChargeResult replay(StoredIdempotency existing, String requestHash, String key) { diff --git a/src/main/java/com/lmoraesdev/payment/application/usecase/ProcessWebhookService.java b/src/main/java/com/lmoraesdev/payment/application/usecase/ProcessWebhookService.java index 1f837a2..21c7d5c 100644 --- a/src/main/java/com/lmoraesdev/payment/application/usecase/ProcessWebhookService.java +++ b/src/main/java/com/lmoraesdev/payment/application/usecase/ProcessWebhookService.java @@ -8,6 +8,7 @@ import com.lmoraesdev.payment.config.logging.Logger5w1hBuilder; import com.lmoraesdev.payment.domain.event.ChargeStatusChangedEvent; import com.lmoraesdev.payment.domain.exception.ChargeNotFoundException; +import com.lmoraesdev.payment.domain.exception.InvalidChargeStatusException; import com.lmoraesdev.payment.domain.model.Charge; import com.lmoraesdev.payment.domain.model.ChargeStatus; import io.micrometer.core.instrument.Counter; @@ -56,7 +57,7 @@ public void process(ProcessWebhookCommand command) { .orElseThrow(() -> new ChargeNotFoundException(command.chargeId())); ChargeStatus previousStatus = charge.getStatus(); - ChargeStatus newStatus = ChargeStatus.valueOf(command.status()); + ChargeStatus newStatus = parseStatus(command.status()); charge.transitionTo(newStatus); @@ -89,4 +90,12 @@ public void process(ProcessWebhookCommand command) { .how("processWebhook") .info(); } + + private ChargeStatus parseStatus(String status) { + try { + return ChargeStatus.valueOf(status); + } catch (IllegalArgumentException e) { + throw new InvalidChargeStatusException(status); + } + } } diff --git a/src/main/java/com/lmoraesdev/payment/domain/exception/InvalidChargeStatusException.java b/src/main/java/com/lmoraesdev/payment/domain/exception/InvalidChargeStatusException.java new file mode 100644 index 0000000..4b34a00 --- /dev/null +++ b/src/main/java/com/lmoraesdev/payment/domain/exception/InvalidChargeStatusException.java @@ -0,0 +1,8 @@ +package com.lmoraesdev.payment.domain.exception; + +public class InvalidChargeStatusException extends DomainException { + + public InvalidChargeStatusException(String status) { + super("INVALID_CHARGE_STATUS", "Status de charge inválido: %s".formatted(status)); + } +} diff --git a/src/main/java/com/lmoraesdev/payment/domain/model/Charge.java b/src/main/java/com/lmoraesdev/payment/domain/model/Charge.java index ce1c826..a37bf12 100644 --- a/src/main/java/com/lmoraesdev/payment/domain/model/Charge.java +++ b/src/main/java/com/lmoraesdev/payment/domain/model/Charge.java @@ -22,14 +22,21 @@ public class Charge { private ChargeStatus status; private final Instant createdAt; private final Instant expiresAt; + private final Long version; private Charge( - UUID id, Money amount, ChargeStatus status, Instant createdAt, Instant expiresAt) { + UUID id, + Money amount, + ChargeStatus status, + Instant createdAt, + Instant expiresAt, + Long version) { this.id = id; this.amount = amount; this.status = status; this.createdAt = createdAt; this.expiresAt = expiresAt; + this.version = version; } public static Charge create(Money amount) { @@ -41,18 +48,24 @@ public static Charge create(Money amount) { amount, ChargeStatus.ACTIVE, createdAt, - createdAt.plus(EXPIRATION)); + createdAt.plus(EXPIRATION), + null); } public static Charge restore( - UUID id, Money amount, ChargeStatus status, Instant createdAt, Instant expiresAt) { + UUID id, + Money amount, + ChargeStatus status, + Instant createdAt, + Instant expiresAt, + Long version) { Objects.requireNonNull(id, "O id é obrigatório"); Objects.requireNonNull(amount, "O montante (Money) é obrigatório"); Objects.requireNonNull(status, "O status é obrigatório"); Objects.requireNonNull(createdAt, "A data de criação é obrigatória"); Objects.requireNonNull(expiresAt, "A data de expiração é obrigatória"); - return new Charge(id, amount, status, createdAt, expiresAt); + return new Charge(id, amount, status, createdAt, expiresAt, version); } public void transitionTo(ChargeStatus newStatus) { @@ -83,6 +96,10 @@ public Instant getExpiresAt() { return expiresAt; } + public Long getVersion() { + return version; + } + @Override public boolean equals(Object o) { if (this == o) { diff --git a/src/main/resources/db/migration/V5__charges_add_version.sql b/src/main/resources/db/migration/V5__charges_add_version.sql new file mode 100644 index 0000000..c0ccee2 --- /dev/null +++ b/src/main/resources/db/migration/V5__charges_add_version.sql @@ -0,0 +1 @@ +alter table charges add column version bigint not null default 0; diff --git a/src/test/java/com/lmoraesdev/payment/adapter/in/web/GlobalExceptionHandlerTest.java b/src/test/java/com/lmoraesdev/payment/adapter/in/web/GlobalExceptionHandlerTest.java new file mode 100644 index 0000000..f37f245 --- /dev/null +++ b/src/test/java/com/lmoraesdev/payment/adapter/in/web/GlobalExceptionHandlerTest.java @@ -0,0 +1,24 @@ +package com.lmoraesdev.payment.adapter.in.web; + +import static org.assertj.core.api.Assertions.assertThat; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.springframework.dao.OptimisticLockingFailureException; +import org.springframework.http.HttpStatus; +import org.springframework.http.ProblemDetail; + +@DisplayName("GlobalExceptionHandler") +class GlobalExceptionHandlerTest { + + private final GlobalExceptionHandler handler = new GlobalExceptionHandler(); + + @Test + @DisplayName("OptimisticLockingFailureException vira 409 Conflict") + void mapsOptimisticLockingFailureExceptionTo409() { + ProblemDetail problem = + handler.handleOptimisticLock(new OptimisticLockingFailureException("stale row")); + + assertThat(problem.getStatus()).isEqualTo(HttpStatus.CONFLICT.value()); + } +} diff --git a/src/test/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeRepositoryIT.java b/src/test/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeRepositoryIT.java index 476a0bf..6c1b121 100644 --- a/src/test/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeRepositoryIT.java +++ b/src/test/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeRepositoryIT.java @@ -37,6 +37,7 @@ void roundTripPreservesAllFields() { assertThat(result.getAmount().amount()).isEqualByComparingTo(charge.getAmount().amount()); assertThat(result.getStatus()).isEqualTo(charge.getStatus()); assertThat(result.getCreatedAt()).isEqualTo(charge.getCreatedAt()); + assertThat(result.getVersion()).isEqualTo(0L); } @Test @@ -58,7 +59,8 @@ void findExpiredActiveReturnsOverdueActiveCharges() { ChargeTestData.money("10.00"), ChargeStatus.ACTIVE, createdAt, - overdueExpiresAt); + overdueExpiresAt, + null); chargeRepository.save(overdue); Charge notYetExpired = ChargeTestData.aCharge().withStatus(ChargeStatus.ACTIVE).build(); diff --git a/src/test/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationCoordinatorTest.java b/src/test/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationCoordinatorTest.java new file mode 100644 index 0000000..bbdd3e8 --- /dev/null +++ b/src/test/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationCoordinatorTest.java @@ -0,0 +1,48 @@ +package com.lmoraesdev.payment.adapter.out.scheduling; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.lmoraesdev.payment.application.port.out.ChargeRepository; +import com.lmoraesdev.payment.application.port.out.OutboxEventPort; +import com.lmoraesdev.payment.domain.model.Charge; +import com.lmoraesdev.payment.domain.model.ChargeStatus; +import com.lmoraesdev.payment.testdata.ChargeTestData; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentMatchers; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +@DisplayName("ChargeExpirationCoordinator") +@ExtendWith(MockitoExtension.class) +class ChargeExpirationCoordinatorTest { + + @Mock ChargeRepository chargeRepository; + + @Mock OutboxEventPort outboxEventPort; + + @InjectMocks ChargeExpirationCoordinator coordinator; + + @Test + @DisplayName("expireOne transiciona pra EXPIRED, salva e grava outbox") + void expireOneTransitionsSavesAndRecordsOutbox() { + Charge charge = ChargeTestData.aCharge().withStatus(ChargeStatus.ACTIVE).build(); + when(chargeRepository.save(charge)).thenReturn(charge); + + coordinator.expireOne(charge); + + assertThat(charge.getStatus()).isEqualTo(ChargeStatus.EXPIRED); + verify(chargeRepository).save(charge); + verify(outboxEventPort) + .record( + eq("Charge"), + eq(charge.getId().toString()), + eq("ChargeExpired"), + ArgumentMatchers.any()); + } +} diff --git a/src/test/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationJobTest.java b/src/test/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationJobTest.java index aba5903..ac0917d 100644 --- a/src/test/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationJobTest.java +++ b/src/test/java/com/lmoraesdev/payment/adapter/out/scheduling/ChargeExpirationJobTest.java @@ -1,14 +1,13 @@ package com.lmoraesdev.payment.adapter.out.scheduling; -import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatCode; import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import com.lmoraesdev.payment.application.port.out.ChargeRepository; -import com.lmoraesdev.payment.application.port.out.OutboxEventPort; import com.lmoraesdev.payment.domain.model.Charge; import com.lmoraesdev.payment.domain.model.ChargeStatus; import com.lmoraesdev.payment.testdata.ChargeTestData; @@ -19,6 +18,7 @@ import org.mockito.InjectMocks; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.dao.OptimisticLockingFailureException; @DisplayName("ChargeExpirationJob") @ExtendWith(MockitoExtension.class) @@ -26,32 +26,46 @@ class ChargeExpirationJobTest { @Mock ChargeRepository chargeRepository; - @Mock OutboxEventPort outboxEventPort; + @Mock ChargeExpirationCoordinator chargeExpirationCoordinator; @InjectMocks ChargeExpirationJob job; @Test - @DisplayName("charge vencida transiciona pra EXPIRED, salva e grava outbox") - void expiresOverdueCharge() { + @DisplayName("charge vencida é reivindicada via ChargeExpirationCoordinator") + void expiresOverdueChargeThroughCoordinator() { Charge charge = ChargeTestData.aCharge().withStatus(ChargeStatus.ACTIVE).build(); when(chargeRepository.findExpiredActive(any())).thenReturn(List.of(charge)); job.expireOverdueCharges(); - assertThat(charge.getStatus()).isEqualTo(ChargeStatus.EXPIRED); - verify(chargeRepository).save(charge); - verify(outboxEventPort) - .record(eq("Charge"), eq(charge.getId().toString()), eq("ChargeExpired"), any()); + verify(chargeExpirationCoordinator).expireOne(charge); } @Test - @DisplayName("lista vazia não faz nada") + @DisplayName("lista vazia não chama o coordinator") void doesNothingWhenListIsEmpty() { when(chargeRepository.findExpiredActive(any())).thenReturn(List.of()); job.expireOverdueCharges(); - verify(chargeRepository, never()).save(any()); - verify(outboxEventPort, never()).record(any(), any(), any(), any()); + verify(chargeExpirationCoordinator, never()).expireOne(any()); + } + + @Test + @DisplayName( + "conflito de otimistic locking numa charge não impede as outras de serem processadas" + + " nem propaga") + void skipsChargeOnOptimisticLockConflictWithoutStoppingOthers() { + Charge conflicting = ChargeTestData.aCharge().withStatus(ChargeStatus.ACTIVE).build(); + Charge healthy = ChargeTestData.aCharge().withStatus(ChargeStatus.ACTIVE).build(); + when(chargeRepository.findExpiredActive(any())).thenReturn(List.of(conflicting, healthy)); + doThrow(new OptimisticLockingFailureException("stale charge")) + .when(chargeExpirationCoordinator) + .expireOne(conflicting); + + assertThatCode(() -> job.expireOverdueCharges()).doesNotThrowAnyException(); + + verify(chargeExpirationCoordinator).expireOne(conflicting); + verify(chargeExpirationCoordinator).expireOne(healthy); } } diff --git a/src/test/java/com/lmoraesdev/payment/application/usecase/ChargeCreationCoordinatorTest.java b/src/test/java/com/lmoraesdev/payment/application/usecase/ChargeCreationCoordinatorTest.java new file mode 100644 index 0000000..9410771 --- /dev/null +++ b/src/test/java/com/lmoraesdev/payment/application/usecase/ChargeCreationCoordinatorTest.java @@ -0,0 +1,86 @@ +package com.lmoraesdev.payment.application.usecase; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.lmoraesdev.payment.application.port.in.CreateChargeCommand; +import com.lmoraesdev.payment.application.port.in.CreateChargeResult; +import com.lmoraesdev.payment.application.port.out.ChargeRepository; +import com.lmoraesdev.payment.application.port.out.IdempotencyPort; +import com.lmoraesdev.payment.application.port.out.OutboxEventPort; +import com.lmoraesdev.payment.domain.model.ChargeStatus; +import com.lmoraesdev.payment.testdata.ChargeTestData; +import io.micrometer.core.instrument.simple.SimpleMeterRegistry; +import java.math.BigDecimal; +import java.util.stream.Stream; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.extension.ExtendWith; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.MethodSource; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +@DisplayName("ChargeCreationCoordinator") +@ExtendWith(MockitoExtension.class) +class ChargeCreationCoordinatorTest { + + @Mock ChargeRepository chargeRepository; + + @Mock OutboxEventPort outboxEventPort; + + @Mock IdempotencyPort idempotencyPort; + + SimpleMeterRegistry meterRegistry; + + ChargeCreationCoordinator coordinator; + + @BeforeEach + void setUp() { + meterRegistry = new SimpleMeterRegistry(); + coordinator = + new ChargeCreationCoordinator( + chargeRepository, outboxEventPort, idempotencyPort, meterRegistry); + } + + record Case(String name, String amount) { + @Override + public String toString() { + return name; + } + } + + static Stream validAmounts() { + return Stream.of( + new Case("centavo mínimo", "0.01"), + new Case("valor comum", "100.00"), + new Case("valor alto", "50000.00")); + } + + @ParameterizedTest + @MethodSource("validAmounts") + @DisplayName("cria cobrança nova, grava outbox e registra idempotency record") + void createAndPersistCreatesChargeSuccessfully(Case c) { + when(chargeRepository.save(any())).thenAnswer(inv -> inv.getArgument(0)); + + CreateChargeResult result = + coordinator.createAndPersist( + new CreateChargeCommand(new BigDecimal(c.amount()), "key-" + c.name()), + ChargeTestData.money(c.amount()), + "hash-" + c.name()); + + assertThat(result.id()).isNotNull(); + assertThat(result.status()).isEqualTo(ChargeStatus.ACTIVE.name()); + assertThat(result.amount()).isEqualByComparingTo(c.amount()); + assertThat(result.createdAt()).isNotNull(); + assertThat(result.replayed()).isFalse(); + verify(chargeRepository).save(any()); + verify(outboxEventPort).record(eq("Charge"), any(), eq("ChargeCreated"), any()); + verify(idempotencyPort) + .save(eq("key-" + c.name()), eq("hash-" + c.name()), any(), eq(result)); + assertThat(meterRegistry.counter("charges_created_total").count()).isEqualTo(1.0); + } +} diff --git a/src/test/java/com/lmoraesdev/payment/application/usecase/CreateChargeServiceConcurrencyIT.java b/src/test/java/com/lmoraesdev/payment/application/usecase/CreateChargeServiceConcurrencyIT.java new file mode 100644 index 0000000..6f9e5ec --- /dev/null +++ b/src/test/java/com/lmoraesdev/payment/application/usecase/CreateChargeServiceConcurrencyIT.java @@ -0,0 +1,71 @@ +package com.lmoraesdev.payment.application.usecase; + +import static org.assertj.core.api.Assertions.assertThat; + +import com.lmoraesdev.payment.adapter.out.persistence.SpringDataChargeRepository; +import com.lmoraesdev.payment.application.port.in.CreateChargeCommand; +import com.lmoraesdev.payment.application.port.in.CreateChargeResult; +import com.lmoraesdev.payment.support.TestcontainersConfiguration; +import java.math.BigDecimal; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.stream.IntStream; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.context.annotation.Import; +import org.springframework.test.context.ActiveProfiles; + +@SpringBootTest +@Import(TestcontainersConfiguration.class) +@ActiveProfiles("test") +@DisplayName("CreateChargeService (concorrência real via Postgres)") +class CreateChargeServiceConcurrencyIT { + + @Autowired CreateChargeService createChargeService; + + @Autowired SpringDataChargeRepository springDataChargeRepository; + + @Test + @DisplayName( + "duas requisições concorrentes com a mesma Idempotency-Key nova: exatamente uma" + + " charge criada, a perdedora recebe o replay da vencedora") + void concurrentRequestsWithSameNewIdempotencyKeyResultInSingleChargeAndReplay() + throws Exception { + String idempotencyKey = UUID.randomUUID().toString(); + CreateChargeCommand command = + new CreateChargeCommand(new BigDecimal("100.00"), idempotencyKey); + CyclicBarrier barrier = new CyclicBarrier(2); + ExecutorService executor = Executors.newFixedThreadPool(2); + + try { + List> futures = + IntStream.range(0, 2) + .mapToObj( + i -> + executor.submit( + () -> { + barrier.await(); + return createChargeService.create(command); + })) + .toList(); + + CreateChargeResult first = futures.get(0).get(10, TimeUnit.SECONDS); + CreateChargeResult second = futures.get(1).get(10, TimeUnit.SECONDS); + + assertThat(first.id()).isEqualTo(second.id()); + assertThat(first.replayed() ^ second.replayed()) + .as("exatamente uma das duas respostas deve ser um replay") + .isTrue(); + assertThat(springDataChargeRepository.count()).isEqualTo(1); + } finally { + executor.shutdownNow(); + } + } +} diff --git a/src/test/java/com/lmoraesdev/payment/application/usecase/CreateChargeServiceTest.java b/src/test/java/com/lmoraesdev/payment/application/usecase/CreateChargeServiceTest.java index 41d4384..b2ae8cc 100644 --- a/src/test/java/com/lmoraesdev/payment/application/usecase/CreateChargeServiceTest.java +++ b/src/test/java/com/lmoraesdev/payment/application/usecase/CreateChargeServiceTest.java @@ -3,21 +3,16 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import com.lmoraesdev.payment.application.port.in.CreateChargeCommand; import com.lmoraesdev.payment.application.port.in.CreateChargeResult; -import com.lmoraesdev.payment.application.port.out.ChargeRepository; import com.lmoraesdev.payment.application.port.out.IdempotencyPort; import com.lmoraesdev.payment.application.port.out.IdempotencyPort.StoredIdempotency; -import com.lmoraesdev.payment.application.port.out.OutboxEventPort; import com.lmoraesdev.payment.domain.exception.IdempotencyConflictException; import com.lmoraesdev.payment.domain.exception.InvalidAmountException; -import com.lmoraesdev.payment.domain.model.ChargeStatus; -import io.micrometer.core.instrument.simple.SimpleMeterRegistry; import java.math.BigDecimal; import java.nio.charset.StandardCharsets; import java.security.MessageDigest; @@ -35,27 +30,21 @@ import org.junit.jupiter.params.provider.MethodSource; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.dao.DataIntegrityViolationException; @DisplayName("CreateChargeService") @ExtendWith(MockitoExtension.class) class CreateChargeServiceTest { - @Mock ChargeRepository chargeRepository; - - @Mock OutboxEventPort outboxEventPort; + @Mock ChargeCreationCoordinator chargeCreationCoordinator; @Mock IdempotencyPort idempotencyPort; - SimpleMeterRegistry meterRegistry; - CreateChargeService service; @BeforeEach void setUp() { - meterRegistry = new SimpleMeterRegistry(); - service = - new CreateChargeService( - chargeRepository, outboxEventPort, idempotencyPort, meterRegistry); + service = new CreateChargeService(chargeCreationCoordinator, idempotencyPort); } record Case(String name, String amount) { @@ -82,35 +71,38 @@ private static String sha256(String value) { } } + private static CreateChargeResult aResult(String amount, boolean replayed) { + return new CreateChargeResult( + UUID.randomUUID(), "ACTIVE", new BigDecimal(amount), Instant.now(), replayed); + } + @ParameterizedTest @MethodSource("validAmounts") - @DisplayName("cria cobrança nova, grava outbox e registra idempotency record") - void createsChargeSuccessfully(Case c) { - when(chargeRepository.save(any())).thenAnswer(inv -> inv.getArgument(0)); + @DisplayName( + "sem registro de idempotência existente, delega criação ao ChargeCreationCoordinator") + void delegatesCreationToCoordinatorWhenNoExistingIdempotencyRecord(Case c) { + when(idempotencyPort.findByKey("key-" + c.name())).thenReturn(Optional.empty()); + CreateChargeResult expected = aResult(c.amount(), false); + when(chargeCreationCoordinator.createAndPersist(any(), any(), any())).thenReturn(expected); CreateChargeResult result = service.create( new CreateChargeCommand(new BigDecimal(c.amount()), "key-" + c.name())); - assertThat(result.id()).isNotNull(); - assertThat(result.status()).isEqualTo(ChargeStatus.ACTIVE.name()); - assertThat(result.amount()).isEqualByComparingTo(c.amount()); - assertThat(result.createdAt()).isNotNull(); - assertThat(result.replayed()).isFalse(); - verify(chargeRepository).save(any()); - verify(outboxEventPort).record(eq("Charge"), any(), eq("ChargeCreated"), any()); - verify(idempotencyPort).save(eq("key-" + c.name()), any(), any(), eq(result)); - assertThat(meterRegistry.counter("charges_created_total").count()).isEqualTo(1.0); + assertThat(result).isEqualTo(expected); + verify(chargeCreationCoordinator).createAndPersist(any(), any(), any()); } @Test - @DisplayName("propaga InvalidAmountException para amount zero ou negativo") + @DisplayName("propaga InvalidAmountException para amount zero ou negativo sem tocar em nada") void propagatesExceptionForInvalidAmount() { assertThatThrownBy( () -> service.create( new CreateChargeCommand(BigDecimal.ZERO, "key-invalid"))) .isInstanceOf(InvalidAmountException.class); + verify(idempotencyPort, never()).findByKey(any()); + verify(chargeCreationCoordinator, never()).createAndPersist(any(), any(), any()); } @Test @@ -121,15 +113,9 @@ void propagatesExceptionForNullAmount() { } @Test - @DisplayName("idempotency key repetida com mesmo body retorna replay sem criar charge nova") + @DisplayName("idempotency key repetida com mesmo body retorna replay sem chamar coordinator") void returnsReplayForRepeatedKeyWithSameBody() { - CreateChargeResult stored = - new CreateChargeResult( - UUID.randomUUID(), - "ACTIVE", - new BigDecimal("100.00"), - Instant.now(), - false); + CreateChargeResult stored = aResult("100.00", false); when(idempotencyPort.findByKey("key-replay")) .thenReturn(Optional.of(new StoredIdempotency(sha256("100.00"), stored))); @@ -141,21 +127,13 @@ void returnsReplayForRepeatedKeyWithSameBody() { assertThat(result.amount()).isEqualByComparingTo(stored.amount()); assertThat(result.createdAt()).isEqualTo(stored.createdAt()); assertThat(result.replayed()).isTrue(); - verify(chargeRepository, never()).save(any()); - verify(outboxEventPort, never()).record(any(), any(), any(), any()); - assertThat(meterRegistry.counter("charges_created_total").count()).isEqualTo(0.0); + verify(chargeCreationCoordinator, never()).createAndPersist(any(), any(), any()); } @Test @DisplayName("idempotency key repetida com body diferente lança IdempotencyConflictException") void throwsConflictForRepeatedKeyWithDifferentBody() { - CreateChargeResult stored = - new CreateChargeResult( - UUID.randomUUID(), - "ACTIVE", - new BigDecimal("100.00"), - Instant.now(), - false); + CreateChargeResult stored = aResult("100.00", false); when(idempotencyPort.findByKey("key-conflict")) .thenReturn(Optional.of(new StoredIdempotency(sha256("100.00"), stored))); @@ -165,6 +143,43 @@ void throwsConflictForRepeatedKeyWithDifferentBody() { new CreateChargeCommand( new BigDecimal("200.00"), "key-conflict"))) .isInstanceOf(IdempotencyConflictException.class); - verify(chargeRepository, never()).save(any()); + verify(chargeCreationCoordinator, never()).createAndPersist(any(), any(), any()); + } + + @Test + @DisplayName( + "duas requisições concorrentes com a mesma key nova: violação de constraint no" + + " coordinator resulta em replay da vencedora, não em erro cru") + void returnsWinnersReplayWhenCoordinatorThrowsConstraintViolation() { + CreateChargeResult winner = aResult("100.00", false); + when(idempotencyPort.findByKey("key-race")) + .thenReturn(Optional.empty()) + .thenReturn(Optional.of(new StoredIdempotency(sha256("100.00"), winner))); + when(chargeCreationCoordinator.createAndPersist(any(), any(), any())) + .thenThrow(new DataIntegrityViolationException("duplicate key")); + + CreateChargeResult result = + service.create(new CreateChargeCommand(new BigDecimal("100.00"), "key-race")); + + assertThat(result.id()).isEqualTo(winner.id()); + assertThat(result.replayed()).isTrue(); + verify(idempotencyPort, org.mockito.Mockito.times(2)).findByKey("key-race"); + } + + @Test + @DisplayName( + "violação de constraint sem registro de idempotência encontrado depois propaga" + + " IllegalStateException") + void propagatesIllegalStateExceptionWhenNoRecordFoundAfterConstraintViolation() { + when(idempotencyPort.findByKey("key-anomaly")).thenReturn(Optional.empty()); + when(chargeCreationCoordinator.createAndPersist(any(), any(), any())) + .thenThrow(new DataIntegrityViolationException("duplicate key")); + + assertThatThrownBy( + () -> + service.create( + new CreateChargeCommand( + new BigDecimal("100.00"), "key-anomaly"))) + .isInstanceOf(IllegalStateException.class); } } diff --git a/src/test/java/com/lmoraesdev/payment/application/usecase/ProcessWebhookServiceTest.java b/src/test/java/com/lmoraesdev/payment/application/usecase/ProcessWebhookServiceTest.java index 99ea4c6..d5baa2d 100644 --- a/src/test/java/com/lmoraesdev/payment/application/usecase/ProcessWebhookServiceTest.java +++ b/src/test/java/com/lmoraesdev/payment/application/usecase/ProcessWebhookServiceTest.java @@ -13,6 +13,7 @@ import com.lmoraesdev.payment.application.port.out.OutboxEventPort; import com.lmoraesdev.payment.application.port.out.WebhookEventPort; import com.lmoraesdev.payment.domain.exception.ChargeNotFoundException; +import com.lmoraesdev.payment.domain.exception.InvalidChargeStatusException; import com.lmoraesdev.payment.domain.exception.InvalidStateTransitionException; import com.lmoraesdev.payment.domain.model.Charge; import com.lmoraesdev.payment.domain.model.ChargeStatus; @@ -26,6 +27,7 @@ import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.dao.OptimisticLockingFailureException; @DisplayName("ProcessWebhookService") @ExtendWith(MockitoExtension.class) @@ -125,4 +127,43 @@ void throwsInvalidStateTransitionExceptionForInvalidTransition() { verify(chargeRepository, never()).save(any()); verify(webhookEventPort, never()).save(any(), any()); } + + @Test + @DisplayName( + "status recebido inválido lança InvalidChargeStatusException antes de transicionar") + void throwsInvalidChargeStatusExceptionForUnknownStatus() { + Charge charge = ChargeTestData.aCharge().withStatus(ChargeStatus.ACTIVE).build(); + when(webhookEventPort.existsByEventId("event-6")).thenReturn(false); + when(chargeRepository.findById(charge.getId())).thenReturn(Optional.of(charge)); + + assertThatThrownBy( + () -> + service.process( + new ProcessWebhookCommand( + "event-6", charge.getId(), "BOGUS_STATUS"))) + .isInstanceOf(InvalidChargeStatusException.class); + assertThat(charge.getStatus()).isEqualTo(ChargeStatus.ACTIVE); + verify(chargeRepository, never()).save(any()); + verify(outboxEventPort, never()).record(any(), any(), any(), any()); + verify(webhookEventPort, never()).save(any(), any()); + } + + @Test + @DisplayName( + "conflito de otimistic locking ao salvar propaga OptimisticLockingFailureException") + void propagatesOptimisticLockingFailureExceptionOnConcurrentUpdate() { + Charge charge = ChargeTestData.aCharge().withStatus(ChargeStatus.ACTIVE).build(); + when(webhookEventPort.existsByEventId("event-7")).thenReturn(false); + when(chargeRepository.findById(charge.getId())).thenReturn(Optional.of(charge)); + when(chargeRepository.save(charge)) + .thenThrow(new OptimisticLockingFailureException("stale charge")); + + assertThatThrownBy( + () -> + service.process( + new ProcessWebhookCommand( + "event-7", charge.getId(), "PAID"))) + .isInstanceOf(OptimisticLockingFailureException.class); + verify(webhookEventPort, never()).save(any(), any()); + } } diff --git a/src/test/java/com/lmoraesdev/payment/domain/model/ChargeTest.java b/src/test/java/com/lmoraesdev/payment/domain/model/ChargeTest.java index 81f00b0..f6ef1db 100644 --- a/src/test/java/com/lmoraesdev/payment/domain/model/ChargeTest.java +++ b/src/test/java/com/lmoraesdev/payment/domain/model/ChargeTest.java @@ -50,13 +50,14 @@ void restorePreservesIdAndTimestamp() { Instant createdAt = Instant.parse("2025-01-01T00:00:00Z"); Instant expiresAt = createdAt.plus(Duration.ofMinutes(30)); - Charge charge = Charge.restore(id, amount, ChargeStatus.PAID, createdAt, expiresAt); + Charge charge = Charge.restore(id, amount, ChargeStatus.PAID, createdAt, expiresAt, 3L); assertThat(charge.getId()).isEqualTo(id); assertThat(charge.getStatus()).isEqualTo(ChargeStatus.PAID); assertThat(charge.getAmount()).isEqualTo(amount); assertThat(charge.getCreatedAt()).isEqualTo(createdAt); assertThat(charge.getExpiresAt()).isEqualTo(expiresAt); + assertThat(charge.getVersion()).isEqualTo(3L); } @Test @@ -66,21 +67,23 @@ void equalityIsIdBased() { Instant now = Instant.now(); Charge a = Charge.restore( - id, amount, ChargeStatus.ACTIVE, now, now.plus(Duration.ofMinutes(30))); + id, amount, ChargeStatus.ACTIVE, now, now.plus(Duration.ofMinutes(30)), 0L); Charge b = Charge.restore( id, new Money(new BigDecimal("99.00")), ChargeStatus.PAID, now, - now.plus(Duration.ofMinutes(30))); + now.plus(Duration.ofMinutes(30)), + 1L); Charge c = Charge.restore( UUID.randomUUID(), amount, ChargeStatus.ACTIVE, now, - now.plus(Duration.ofMinutes(30))); + now.plus(Duration.ofMinutes(30)), + 0L); assertThat(a).isEqualTo(b); assertThat(a).isNotEqualTo(c); @@ -104,7 +107,7 @@ void transitionToRejectsInvalidTransition() { Instant now = Instant.now(); Charge charge = Charge.restore( - id, amount, ChargeStatus.PAID, now, now.plus(Duration.ofMinutes(30))); + id, amount, ChargeStatus.PAID, now, now.plus(Duration.ofMinutes(30)), 0L); assertThatThrownBy(() -> charge.transitionTo(ChargeStatus.ACTIVE)) .isInstanceOf(InvalidStateTransitionException.class) diff --git a/src/test/java/com/lmoraesdev/payment/testdata/ChargeTestData.java b/src/test/java/com/lmoraesdev/payment/testdata/ChargeTestData.java index 322b204..f1a3eab 100644 --- a/src/test/java/com/lmoraesdev/payment/testdata/ChargeTestData.java +++ b/src/test/java/com/lmoraesdev/payment/testdata/ChargeTestData.java @@ -12,6 +12,7 @@ public final class ChargeTestData { private Money amount = new Money(new BigDecimal("10.50")); private ChargeStatus status = ChargeStatus.ACTIVE; + private Long version = 0L; private ChargeTestData() {} @@ -33,9 +34,14 @@ public ChargeTestData withStatus(ChargeStatus s) { return this; } + public ChargeTestData withVersion(Long v) { + this.version = v; + return this; + } + public Charge build() { Instant now = Instant.now(); return Charge.restore( - UUID.randomUUID(), amount, status, now, now.plus(Duration.ofMinutes(30))); + UUID.randomUUID(), amount, status, now, now.plus(Duration.ofMinutes(30)), version); } } From 485bc7f85c0c78af79c46f832bc6114e901c8535 Mon Sep 17 00:00:00 2001 From: Leandro Moraes Date: Thu, 23 Jul 2026 13:48:21 -0300 Subject: [PATCH 3/4] fix(docs): corrige URL do badge de CI no README --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 59b633c..e21cbde 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ Core de pagamentos Pix em Java — Hexagonal Architecture, observabilidade e pipeline CI/CD. EPIC-001 (criar cobrança) implementado e testado. -[![CI](https://github.com/lmoraesdev/java-payment-hexagonal/actions/workflows/ci.yml/badge.svg)](https://github.com/lmoraesdev/java-payment-hexagonal/actions/workflows/ci.yml) +[![CI](https://github.com/lmoraesdev/java-payment-core/actions/workflows/ci.yml/badge.svg)](https://github.com/lmoraesdev/java-payment-core/actions/workflows/ci.yml) ![Java](https://img.shields.io/badge/Java-21-blue?logo=openjdk&logoColor=white) ![Spring Boot](https://img.shields.io/badge/Spring%20Boot-3.5.3-6DB33F?logo=springboot&logoColor=white) From d6005ba754c99531fbf0fd0b55f17d2408cd347f Mon Sep 17 00:00:00 2001 From: Leandro Moraes Date: Thu, 23 Jul 2026 13:50:17 -0300 Subject: [PATCH 4/4] fix(charges): corrige falso optimistic lock em novas charges e ajusta testes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ChargeJpaEntity agora implementa Persistable com isNew() baseado em (version == null). O id da Charge é atribuído em código (UUID.randomUUID()), não pelo banco, então Spring Data não tinha como distinguir insert de update só pelo id e sempre chamava merge() — com @Version presente, merge() de uma entidade nova faz o Hibernate suspeitar que a linha foi apagada por outra transação e lança ObjectOptimisticLockingFailureException. Corrige também OutboxClaimCoordinatorTransactionalIT: removido um thenCallRealMethod() estubado num método de repository Spring Data (proxy sem corpo Java real, Mockito não consegue chamar); o @MockitoSpyBean já delega pro objeto real automaticamente em qualquer chamada não estubada. Corrige ChargeRepositoryIT: os dois casos que esperavam inserir uma charge nova usavam ChargeTestData.aCharge().build(), que via Charge.restore() monta uma charge com version=0 (representando uma linha já persistida). Isso fazia isNew() reportar corretamente false e o save() tentar um merge() numa linha inexistente. Trocado por Charge.create(...), a fábrica de domínio real usada em produção para charges novas (version=null). --- .../adapter/out/persistence/ChargeJpaEntity.java | 9 ++++++++- .../OutboxClaimCoordinatorTransactionalIT.java | 10 ++-------- .../adapter/out/persistence/ChargeRepositoryIT.java | 4 ++-- 3 files changed, 12 insertions(+), 11 deletions(-) diff --git a/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeJpaEntity.java b/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeJpaEntity.java index 59d3b4e..c3050b6 100644 --- a/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeJpaEntity.java +++ b/src/main/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeJpaEntity.java @@ -10,10 +10,11 @@ import jakarta.persistence.Version; import java.time.Instant; import java.util.UUID; +import org.springframework.data.domain.Persistable; @Entity @Table(name = "charges") -public class ChargeJpaEntity { +public class ChargeJpaEntity implements Persistable { @Id private UUID id; @Column(name = "amount_centavos", nullable = false) @@ -35,10 +36,16 @@ public class ChargeJpaEntity { protected ChargeJpaEntity() {} + @Override public UUID getId() { return id; } + @Override + public boolean isNew() { + return version == null; + } + public void setId(UUID id) { this.id = id; } diff --git a/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTransactionalIT.java b/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTransactionalIT.java index b07ce96..97f562a 100644 --- a/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTransactionalIT.java +++ b/src/test/java/com/lmoraesdev/payment/adapter/out/messaging/OutboxClaimCoordinatorTransactionalIT.java @@ -3,7 +3,7 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.anyList; -import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doThrow; import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventEntity; import com.lmoraesdev.payment.adapter.out.persistence.outbox.OutboxEventJpaRepository; @@ -46,13 +46,7 @@ void exceptionMidClaimBatchRollsBackStatusChange() { OutboxEventEntity.pending( "Charge", "aggregate-rollback", "ChargeCreated", "{}", null)); - doAnswer( - invocation -> { - invocation.callRealMethod(); - throw new RuntimeException("boom-mid-claim"); - }) - .when(repository) - .saveAll(anyList()); + doThrow(new RuntimeException("boom-mid-claim")).when(repository).saveAll(anyList()); assertThatThrownBy(() -> coordinator.claimBatch()) .isInstanceOf(RuntimeException.class) diff --git a/src/test/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeRepositoryIT.java b/src/test/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeRepositoryIT.java index 6c1b121..6364eae 100644 --- a/src/test/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeRepositoryIT.java +++ b/src/test/java/com/lmoraesdev/payment/adapter/out/persistence/ChargeRepositoryIT.java @@ -26,7 +26,7 @@ class ChargeRepositoryIT extends AbstractIntegrationTest { @Test @DisplayName("save e findById preservam todos os campos") void roundTripPreservesAllFields() { - Charge charge = ChargeTestData.aCharge().build(); + Charge charge = Charge.create(ChargeTestData.money("10.50")); Charge saved = chargeRepository.save(charge); Optional found = chargeRepository.findById(saved.getId()); @@ -63,7 +63,7 @@ void findExpiredActiveReturnsOverdueActiveCharges() { null); chargeRepository.save(overdue); - Charge notYetExpired = ChargeTestData.aCharge().withStatus(ChargeStatus.ACTIVE).build(); + Charge notYetExpired = Charge.create(ChargeTestData.money("10.00")); chargeRepository.save(notYetExpired); List result = chargeRepository.findExpiredActive(Instant.now());