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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
@@ -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<OutboxEventEntity> claimBatch() {
repository.reapStuckInFlight(Instant.now().minus(STUCK_IN_FLIGHT_THRESHOLD));

List<OutboxEventEntity> 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);
});
}
}
Original file line number Diff line number Diff line change
@@ -1,37 +1,34 @@
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;
import io.micrometer.core.instrument.Timer;
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<Object, Object> kafkaTemplate;
private final Counter publishedCounter;
private final Counter failedCounter;
private final Timer publishLagTimer;

public OutboxRelay(
OutboxEventJpaRepository repository,
OutboxClaimCoordinator outboxClaimCoordinator,
KafkaTemplate<Object, Object> 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");
Expand All @@ -40,30 +37,23 @@ public OutboxRelay(

@Scheduled(fixedDelay = 5000)
public void publishPending() {
List<OutboxEventEntity> claimed = claimBatch();
List<OutboxEventEntity> claimed = outboxClaimCoordinator.claimBatch();

for (OutboxEventEntity event : claimed) {
publish(event);
}
}

@Transactional
public List<OutboxEventEntity> claimBatch() {
List<OutboxEventEntity> 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());
}
}

Expand All @@ -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);
});
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,12 +7,14 @@
import jakarta.persistence.Enumerated;
import jakarta.persistence.Id;
import jakarta.persistence.Table;
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<UUID> {
@Id private UUID id;

@Column(name = "amount_centavos", nullable = false)
Expand All @@ -28,12 +30,22 @@ public class ChargeJpaEntity {
@Column(name = "expires_at", nullable = false)
private Instant expiresAt;

@Version
@Column(nullable = false)
private Long version;

protected ChargeJpaEntity() {}

@Override
public UUID getId() {
return id;
}

@Override
public boolean isNew() {
return version == null;
}

public void setId(UUID id) {
this.id = id;
}
Expand Down Expand Up @@ -69,4 +81,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;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand All @@ -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());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,9 @@ public class OutboxEventEntity {
@Column(name = "correlation_id")
private String correlationId;

@Column(name = "claimed_at")
private Instant claimedAt;

protected OutboxEventEntity() {
// JPA
}
Expand All @@ -70,6 +73,7 @@ public static OutboxEventEntity pending(

public void markInFlight() {
this.status = OutboxStatus.IN_FLIGHT;
this.claimedAt = Instant.now();
}

public void markPublished() {
Expand All @@ -79,6 +83,7 @@ public void markPublished() {

public void revertToPending() {
this.status = OutboxStatus.PENDING;
this.claimedAt = null;
}

public UUID getId() {
Expand Down Expand Up @@ -108,4 +113,8 @@ public Instant getCreatedAt() {
public String getCorrelationId() {
return correlationId;
}

public Instant getClaimedAt() {
return claimedAt;
}
}
Original file line number Diff line number Diff line change
@@ -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<OutboxEventEntity, UUID> {

Expand All @@ -18,4 +21,15 @@ public interface OutboxEventJpaRepository extends JpaRepository<OutboxEventEntit
""",
nativeQuery = true)
List<OutboxEventEntity> 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);
}
Original file line number Diff line number Diff line change
@@ -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();
}
}
Loading
Loading