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
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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(),
Expand Down
Original file line number Diff line number Diff line change
@@ -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);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
}
Original file line number Diff line number Diff line change
@@ -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);
}
Original file line number Diff line number Diff line change
Expand Up @@ -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로 직행하도록 원본 예외 던짐
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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일 수 없습니다.");
}
Expand Down
Original file line number Diff line number Diff line change
@@ -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<ProcessedEvent, UUID> {
}
Original file line number Diff line number Diff line change
@@ -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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,10 @@ public ResponseEntity<ApiResponse<Void>> reserveStock(
return ResponseEntity.ok(ApiResponse.ok(null));
}

/**
* Kafka 비동기 통신(order.stock-restore.requested)으로 대체됨. 향후 삭제 예정
*/
@Deprecated(since = "1.0", forRemoval = true)
@PostMapping("/restore")
public ResponseEntity<ApiResponse<Void>> restoreStock(
@RequestBody @Valid RestoreStockRequest request
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
) {
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,6 @@
public record RestoreStockRequest(
@NotNull UUID optionId,
@Positive int quantity,
UUID reservationId // 멱등키 — TODO 추후 중복 복구 방지 로직에서 활용 예정 (그땐 @NotNull)
@NotNull UUID eventId // 멱등키 필드
) {
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -51,6 +52,10 @@ class StockCommandServiceTest {
@Mock
private InventoryOutboxHelper outboxHelper;

// 멱등성 검증을 위한 Mock 객체 추가
@Mock
private ProcessedEventRepository processedEventRepository;

// 카프카로 전송된 이벤트를 낚아채서 내부 값을 검증하기 위한 Captor
@Captor
private ArgumentCaptor<StockReservedEvent> eventCaptor;
Expand Down Expand Up @@ -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
Expand All @@ -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))
Expand Down