diff --git a/src/main/java/com/michelet/inventory/application/StockCommandService.java b/src/main/java/com/michelet/inventory/application/StockCommandService.java index 76ef241..26607ac 100644 --- a/src/main/java/com/michelet/inventory/application/StockCommandService.java +++ b/src/main/java/com/michelet/inventory/application/StockCommandService.java @@ -4,10 +4,12 @@ import com.michelet.inventory.application.dto.StockReservedEvent; import com.michelet.inventory.application.dto.StockRestoredEvent; import com.michelet.inventory.domain.exception.StockNotFoundException; +import com.michelet.inventory.domain.model.ProcessedEvent; import com.michelet.inventory.domain.model.Product; import com.michelet.inventory.domain.model.ProductOption; import com.michelet.inventory.domain.model.ProductStatus; import com.michelet.inventory.domain.model.Stock; +import com.michelet.inventory.domain.repository.ProcessedEventRepository; import com.michelet.inventory.domain.repository.ProductOptionRepository; import com.michelet.inventory.domain.repository.ProductRepository; import com.michelet.inventory.domain.repository.StockRepository; @@ -26,6 +28,7 @@ public class StockCommandService { private final StockRepository stockRepository; private final ProductOptionRepository productOptionRepository; private final ProductRepository productRepository; + private final ProcessedEventRepository processedEventRepository; // OutboxHelper 주입 (KafkaTemplate 대체) private final InventoryOutboxHelper outboxHelper; @@ -71,16 +74,25 @@ public void reserveStock(ReserveStockRequest request) { @Transactional public void restoreStock(RestoreStockRequest request) { + // 1. 멱등성 검증 (Redisson Lock 내부이므로 동시성 중복 방어됨) + if (processedEventRepository.existsById(request.eventId())) { + log.warn("[Idempotency] 이미 처리된 카프카 메시지입니다. 중복 복구를 방지하고 반환합니다. eventId={}", request.eventId()); + return; // 중복 메시지는 로직을 건너뛰고 정상 완료 처리 + } + Stock stock = stockRepository.findById(request.optionId()) .orElseThrow(StockNotFoundException::new); - // 1. 도메인 로직: 재고 복구 + // 2. 도메인 로직: 재고 복구 stock.restore(request.quantity()); - // 2. DB 업데이트 (더티 체킹 후 flush) + // 3. DB 업데이트 stockRepository.save(stock); - // 3. 재고 복구 이벤트를 Outbox에 적재 + // 4. 처리 완료 기록 (멱등키 저장) + processedEventRepository.save(new ProcessedEvent(request.eventId())); + + // 5. 재고 복구 결과를 Outbox에 적재하여 카탈로그로 발송 StockRestoredEvent event = new StockRestoredEvent( stock.getOptionId(), stock.getTotalQuantity(), diff --git a/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java b/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java new file mode 100644 index 0000000..48aea8d --- /dev/null +++ b/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java @@ -0,0 +1,32 @@ +package com.michelet.inventory.domain.model; + +import jakarta.persistence.Column; +import jakarta.persistence.Entity; +import jakarta.persistence.Id; +import jakarta.persistence.Table; +import java.time.Instant; +import java.time.temporal.ChronoUnit; +import java.util.Objects; +import java.util.UUID; +import lombok.AccessLevel; +import lombok.Getter; +import lombok.NoArgsConstructor; + +@Entity +@Table(name = "p_processed_events") +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +public class ProcessedEvent { + + @Id + @Column(name = "event_id", updatable = false, nullable = false) + private UUID eventId; + + @Column(nullable = false) + private Instant processedAt; + + public ProcessedEvent(UUID eventId) { + this.eventId = Objects.requireNonNull(eventId, "eventId는 null일 수 없습니다"); + this.processedAt = Instant.now().truncatedTo(ChronoUnit.MILLIS); + } +} diff --git a/src/main/java/com/michelet/inventory/domain/repository/ProcessedEventRepository.java b/src/main/java/com/michelet/inventory/domain/repository/ProcessedEventRepository.java new file mode 100644 index 0000000..0805679 --- /dev/null +++ b/src/main/java/com/michelet/inventory/domain/repository/ProcessedEventRepository.java @@ -0,0 +1,10 @@ +package com.michelet.inventory.domain.repository; + +import com.michelet.inventory.domain.model.ProcessedEvent; +import java.util.UUID; + +public interface ProcessedEventRepository { + boolean existsById(UUID eventId); + + ProcessedEvent save(ProcessedEvent event); +} diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java index fe3bef8..9f9563f 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java @@ -26,15 +26,18 @@ public void consumeStockRestoreRequest(StockRestoreMessage payload) { throw new IllegalArgumentException("재고 복구 이벤트 파싱 오류: payload는 null일 수 없습니다."); } - log.info("[Kafka Consumer] 재고 복구 이벤트 수신: optionId={}, quantity={}", payload.optionId(), payload.quantity()); + log.info("[Kafka Consumer] 재고 복구 이벤트 수신: eventId={}, optionId={}, quantity={}", payload.eventId(), + payload.optionId(), payload.quantity()); try { // 파사드를 통해 분산 락을 걸고 안전하게 재고 복구 로직 실행 - // TODO(`#26`): reservationId를 Kafka 페이로드에 포함시켜 멱등성 키로 활용 - RestoreStockRequest request = new RestoreStockRequest(payload.optionId(), payload.quantity(), null); + // 전달받은 eventId를 파사드 요청에 포함 + RestoreStockRequest request = new RestoreStockRequest(payload.optionId(), payload.quantity(), + payload.eventId()); stockLockFacade.restoreStockWithLock(request); - log.info("[Kafka Consumer] 재고 복구 완료! optionId: {}, quantity: {}", payload.optionId(), payload.quantity()); + log.info("[Kafka Consumer] 재고 복구 로직 처리 완료! eventId: {}, optionId: {}, quantity: {}", payload.eventId(), + payload.optionId(), payload.quantity()); } catch (IllegalArgumentException e) { // 데이터 검증 실패나 비즈니스 룰 위반 시 -> 즉각 DLT로 직행하도록 원본 예외 던짐 diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/StockRestoreMessage.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/StockRestoreMessage.java index 43d3767..25f6b9b 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/StockRestoreMessage.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/StockRestoreMessage.java @@ -4,10 +4,14 @@ // 오더 서비스가 발행한 이벤트를 읽어들이기 위한 수신 전용 DTO public record StockRestoreMessage( + UUID eventId, // 멱등성 검증용 고유 이벤트 ID UUID optionId, Integer quantity ) { public StockRestoreMessage { + if (eventId == null) { + throw new IllegalArgumentException("재고 복구 이벤트 파싱 오류: eventId는 null일 수 없습니다."); + } if (optionId == null) { throw new IllegalArgumentException("재고 복구 이벤트 파싱 오류: optionId는 null일 수 없습니다."); } diff --git a/src/main/java/com/michelet/inventory/infrastructure/repository/JpaProcessedEventRepository.java b/src/main/java/com/michelet/inventory/infrastructure/repository/JpaProcessedEventRepository.java new file mode 100644 index 0000000..cf3a1c7 --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/repository/JpaProcessedEventRepository.java @@ -0,0 +1,8 @@ +package com.michelet.inventory.infrastructure.repository; + +import com.michelet.inventory.domain.model.ProcessedEvent; +import java.util.UUID; +import org.springframework.data.jpa.repository.JpaRepository; + +public interface JpaProcessedEventRepository extends JpaRepository { +} diff --git a/src/main/java/com/michelet/inventory/infrastructure/repository/ProcessedEventRepositoryImpl.java b/src/main/java/com/michelet/inventory/infrastructure/repository/ProcessedEventRepositoryImpl.java new file mode 100644 index 0000000..e56d61f --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/repository/ProcessedEventRepositoryImpl.java @@ -0,0 +1,24 @@ +package com.michelet.inventory.infrastructure.repository; + +import com.michelet.inventory.domain.model.ProcessedEvent; +import com.michelet.inventory.domain.repository.ProcessedEventRepository; +import java.util.UUID; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Repository; + +@Repository +@RequiredArgsConstructor +public class ProcessedEventRepositoryImpl implements ProcessedEventRepository { + + private final JpaProcessedEventRepository jpaRepository; + + @Override + public boolean existsById(UUID eventId) { + return jpaRepository.existsById(eventId); + } + + @Override + public ProcessedEvent save(ProcessedEvent event) { + return jpaRepository.save(event); + } +} diff --git a/src/main/java/com/michelet/inventory/presentation/InternalStockController.java b/src/main/java/com/michelet/inventory/presentation/InternalStockController.java index bdc5c67..030b85d 100644 --- a/src/main/java/com/michelet/inventory/presentation/InternalStockController.java +++ b/src/main/java/com/michelet/inventory/presentation/InternalStockController.java @@ -27,6 +27,10 @@ public ResponseEntity> reserveStock( return ResponseEntity.ok(ApiResponse.ok(null)); } + /** + * Kafka 비동기 통신(order.stock-restore.requested)으로 대체됨. 향후 삭제 예정 + */ + @Deprecated(since = "1.0", forRemoval = true) @PostMapping("/restore") public ResponseEntity> restoreStock( @RequestBody @Valid RestoreStockRequest request diff --git a/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java b/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java index 1d01b4d..8b337fc 100644 --- a/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java +++ b/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java @@ -5,8 +5,12 @@ import java.util.UUID; public record ReserveStockRequest( - @NotNull UUID optionId, - @NotNull @Min(1) Integer quantity, - UUID reservationId // 멱등키 — TODO 추후 중복 차감 방지 로직에서 활용 예정 (그땐 @NotNull) + @NotNull(message = "optionId는 필수입니다.") + UUID optionId, + @NotNull(message = "quantity는 필수입니다.") + @Min(value = 1, message = "quantity는 1 이상이어야 합니다.") + Integer quantity, + @NotNull(message = "멱등키(reservationId)는 필수입니다.") + UUID reservationId ) { } diff --git a/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java b/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java index e3e488e..acb0540 100644 --- a/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java +++ b/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java @@ -7,6 +7,6 @@ public record RestoreStockRequest( @NotNull UUID optionId, @Positive int quantity, - UUID reservationId // 멱등키 — TODO 추후 중복 복구 방지 로직에서 활용 예정 (그땐 @NotNull) + @NotNull UUID eventId // 멱등키 필드 ) { } diff --git a/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java b/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java index 67532c9..b9daac7 100644 --- a/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java +++ b/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java @@ -16,6 +16,7 @@ import com.michelet.inventory.domain.exception.SoldOutException; import com.michelet.inventory.domain.exception.StockNotFoundException; import com.michelet.inventory.domain.model.Stock; +import com.michelet.inventory.domain.repository.ProcessedEventRepository; import com.michelet.inventory.domain.repository.ProductOptionRepository; import com.michelet.inventory.domain.repository.ProductRepository; import com.michelet.inventory.domain.repository.StockRepository; @@ -51,6 +52,10 @@ class StockCommandServiceTest { @Mock private InventoryOutboxHelper outboxHelper; + // 멱등성 검증을 위한 Mock 객체 추가 + @Mock + private ProcessedEventRepository processedEventRepository; + // 카프카로 전송된 이벤트를 낚아채서 내부 값을 검증하기 위한 Captor @Captor private ArgumentCaptor eventCaptor; @@ -139,18 +144,22 @@ void reserveStock_Fail_MaxLimitExceeded() { } @Test - @DisplayName("성공: 재고 복구 경로 정상 동작 및 Outbox 적재 테스트") + @DisplayName("성공: 재고 복구 경로 정상 동작 및 멱등키, Outbox 적재 테스트") void restoreStock_Success() { UUID optionId = UUID.randomUUID(); + UUID eventId = UUID.randomUUID(); // 멱등키(메시지 ID) + // 1. 초기 재고 100개, 일일 재고 50개 생성 Stock stock = Stock.create(optionId, 100, 50, 10); // 2. 소비된 상태를 시뮬레이션하기 위해 미리 2개를 차감 (total: 98, daily: 48) stock.reserve(2); - // 3. 다시 2개를 복구해 달라는 요청 - RestoreStockRequest request = new RestoreStockRequest(optionId, 2, null); + // 3. 다시 2개를 복구해 달라는 요청 - eventId를 포함하여 요청 객체 생성 + RestoreStockRequest request = new RestoreStockRequest(optionId, 2, eventId); + // 처음 들어온 메시지이므로 existsById는 false를 반환하도록 Mocking + given(processedEventRepository.existsById(eventId)).willReturn(false); given(stockRepository.findById(optionId)).willReturn(Optional.of(stock)); // when @@ -169,12 +178,34 @@ void restoreStock_Success() { assertThat(event.currentDailyStock()).isEqualTo(50); } + // 멱등성 보장 핵심 테스트 + @Test + @DisplayName("성공(멱등성): 이미 처리된 이벤트인 경우 이중 복구를 수행하지 않고 무시한다.") + void restoreStock_Idempotency() { + UUID optionId = UUID.randomUUID(); + UUID eventId = UUID.randomUUID(); + RestoreStockRequest request = new RestoreStockRequest(optionId, 2, eventId); + + // DB에 이미 eventId가 저장되어 있다고(처리되었다고) Mocking + given(processedEventRepository.existsById(eventId)).willReturn(true); + + // when + stockCommandService.restoreStock(request); + + // then: DB 조회나 아웃박스 저장이 단 한 번도 호출되지 않아야 함! (이중 복구 방지 증명) + verifyNoInteractions(stockRepository); + verifyNoInteractions(outboxHelper); + } + @Test @DisplayName("실패: 재고 복구 시 존재하지 않는 옵션 예외") void restoreStock_Fail_NotFound() { UUID optionId = UUID.randomUUID(); - RestoreStockRequest request = new RestoreStockRequest(optionId, 5, null); + UUID eventId = UUID.randomUUID(); + RestoreStockRequest request = new RestoreStockRequest(optionId, 5, eventId); + // 멱등성 검사 통과 설정 + given(processedEventRepository.existsById(eventId)).willReturn(false); given(stockRepository.findById(optionId)).willReturn(Optional.empty()); assertThatThrownBy(() -> stockCommandService.restoreStock(request))