diff --git a/src/main/java/com/michelet/inventory/application/InventoryOutboxHelper.java b/src/main/java/com/michelet/inventory/application/InventoryOutboxHelper.java index 1f0e3bd..2dfae00 100644 --- a/src/main/java/com/michelet/inventory/application/InventoryOutboxHelper.java +++ b/src/main/java/com/michelet/inventory/application/InventoryOutboxHelper.java @@ -33,21 +33,10 @@ public void validateMaxRetries() { } } - // MANDATORY로 변경하여 부모 트랜잭션이 없으면 즉각 실패하도록 원자성 강제 + // MANDATORY로 변경하여 부모 트랜잭션이 없으면 즉각 실패하도록 원자성 강제 (정상 비즈니스 로직용) @Transactional(propagation = Propagation.MANDATORY) public void append(String aggregateType, String aggregateId, String eventType, Object payloadObj) { - if (aggregateType == null || aggregateType.isBlank()) { - throw new IllegalArgumentException("aggregateType은 필수입니다."); - } - if (aggregateId == null || aggregateId.isBlank()) { - throw new IllegalArgumentException("aggregateId는 필수입니다."); - } - if (eventType == null || eventType.isBlank()) { - throw new IllegalArgumentException("eventType은 필수입니다."); - } - if (payloadObj == null) { - throw new IllegalArgumentException("payloadObj는 필수입니다."); - } + validateInputs(aggregateType, aggregateId, eventType, payloadObj); try { String payloadJson = objectMapper.writeValueAsString(payloadObj); @@ -65,6 +54,27 @@ public void append(String aggregateType, String aggregateId, String eventType, O } } + // 부모 트랜잭션이 롤백된 후, Catch 블록 등에서 독립적으로 보상/거절 이벤트를 기록할 때 사용 (Saga 롤백용) + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void appendIndependent(String aggregateType, String aggregateId, String eventType, Object payloadObj) { + validateInputs(aggregateType, aggregateId, eventType, payloadObj); + + try { + String payloadJson = objectMapper.writeValueAsString(payloadObj); + InventoryOutbox outbox = InventoryOutbox.builder() + .aggregateType(aggregateType) + .aggregateId(aggregateId) + .eventType(eventType) + .payload(payloadJson) + .build(); + outboxRepository.save(outbox); + log.info("[Inventory Outbox] 독립 트랜잭션 이벤트 적재 요청 (롤백/거절용): type={}, id={}", eventType, aggregateId); + } catch (JsonProcessingException e) { + log.error("Outbox 페이로드 직렬화 실패. aggregateId={}, eventType={}", aggregateId, eventType, e); + throw new RuntimeException("Outbox 이벤트 생성 중 오류가 발생했습니다.", e); + } + } + // 2단계 스케줄러에서 상태 업데이트 시 사용할 독립 트랜잭션 메서드 @Transactional(propagation = Propagation.REQUIRES_NEW) public void markAsPublished(UUID outboxId) { @@ -95,4 +105,20 @@ public void handleFailure(UUID outboxId) { outboxRepository.save(outbox); }); } + + // 공통 파라미터 검증 메서드 + private void validateInputs(String aggregateType, String aggregateId, String eventType, Object payloadObj) { + if (aggregateType == null || aggregateType.isBlank()) { + throw new IllegalArgumentException("aggregateType은 필수입니다."); + } + if (aggregateId == null || aggregateId.isBlank()) { + throw new IllegalArgumentException("aggregateId는 필수입니다."); + } + if (eventType == null || eventType.isBlank()) { + throw new IllegalArgumentException("eventType은 필수입니다."); + } + if (payloadObj == null) { + throw new IllegalArgumentException("payloadObj는 필수입니다."); + } + } } diff --git a/src/main/java/com/michelet/inventory/application/InventoryOutboxScheduler.java b/src/main/java/com/michelet/inventory/application/InventoryOutboxScheduler.java index 4ec3f37..3be604c 100644 --- a/src/main/java/com/michelet/inventory/application/InventoryOutboxScheduler.java +++ b/src/main/java/com/michelet/inventory/application/InventoryOutboxScheduler.java @@ -10,6 +10,8 @@ import com.michelet.inventory.domain.model.InventoryOutbox; import com.michelet.inventory.domain.model.OutboxStatus; import com.michelet.inventory.domain.repository.InventoryOutboxRepository; +import com.michelet.inventory.infrastructure.messaging.dto.OrderApprovedEvent; +import com.michelet.inventory.infrastructure.messaging.dto.OrderRejectedEvent; import java.util.List; import java.util.UUID; import lombok.RequiredArgsConstructor; @@ -39,6 +41,8 @@ public class InventoryOutboxScheduler { private static final String EVENT_STOCK_RESERVED = "STOCK_RESERVED"; private static final String EVENT_STOCK_RESTORED = "STOCK_RESTORED"; private static final String EVENT_DAILY_RESET = "DAILY_STOCK_RESET"; + private static final String EVENT_ORDER_APPROVED = "ORDER_APPROVED"; + private static final String EVENT_ORDER_REJECTED = "ORDER_REJECTED"; @Value("${inventory.kafka.topic.product-created:product.created}") private String topicProductCreated; @@ -52,6 +56,10 @@ public class InventoryOutboxScheduler { private String topicStockRestored; @Value("${inventory.kafka.topic.daily-reset:stock.daily-reset}") private String topicDailyReset; + @Value("${inventory.kafka.topic.order-approved:order.approved}") + private String topicOrderApproved; + @Value("${inventory.kafka.topic.order-rejected:order.rejected}") + private String topicOrderRejected; @Scheduled(fixedDelay = 5000) public void processOutboxEvents() { @@ -119,6 +127,8 @@ private Object deserializePayload(String eventType, String jsonPayload) throws E case EVENT_STOCK_RESERVED -> objectMapper.readValue(jsonPayload, StockReservedEvent.class); case EVENT_STOCK_RESTORED -> objectMapper.readValue(jsonPayload, StockRestoredEvent.class); case EVENT_DAILY_RESET -> objectMapper.readValue(jsonPayload, DailyStockResetEvent.class); + case EVENT_ORDER_APPROVED -> objectMapper.readValue(jsonPayload, OrderApprovedEvent.class); + case EVENT_ORDER_REJECTED -> objectMapper.readValue(jsonPayload, OrderRejectedEvent.class); // 매핑 안 된 이벤트를 String으로 보내면 직렬화 에러 발생! 예외를 던져서 스케줄러 재시도 루프로 넘김 default -> { log.warn("등록되지 않은 알 수 없는 이벤트 타입입니다: {}", eventType); @@ -135,6 +145,8 @@ private String resolveTopic(String eventType) { case EVENT_STOCK_RESERVED -> topicStockReserved; case EVENT_STOCK_RESTORED -> topicStockRestored; case EVENT_DAILY_RESET -> topicDailyReset; + case EVENT_ORDER_APPROVED -> topicOrderApproved; + case EVENT_ORDER_REJECTED -> topicOrderRejected; default -> "inventory.unknown.event"; }; } diff --git a/src/main/java/com/michelet/inventory/application/StockCommandService.java b/src/main/java/com/michelet/inventory/application/StockCommandService.java index 26607ac..4d32f07 100644 --- a/src/main/java/com/michelet/inventory/application/StockCommandService.java +++ b/src/main/java/com/michelet/inventory/application/StockCommandService.java @@ -1,6 +1,8 @@ package com.michelet.inventory.application; import com.michelet.inventory.application.dto.ProductStatusChangedEvent; +import com.michelet.inventory.application.dto.ReserveStockRequest; +import com.michelet.inventory.application.dto.RestoreStockRequest; import com.michelet.inventory.application.dto.StockReservedEvent; import com.michelet.inventory.application.dto.StockRestoredEvent; import com.michelet.inventory.domain.exception.StockNotFoundException; @@ -13,8 +15,13 @@ import com.michelet.inventory.domain.repository.ProductOptionRepository; import com.michelet.inventory.domain.repository.ProductRepository; import com.michelet.inventory.domain.repository.StockRepository; -import com.michelet.inventory.presentation.dto.ReserveStockRequest; -import com.michelet.inventory.presentation.dto.RestoreStockRequest; +import com.michelet.inventory.infrastructure.messaging.dto.OrderApprovedEvent; +import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.stream.Collectors; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -100,4 +107,64 @@ public void restoreStock(RestoreStockRequest request) { ); outboxHelper.append("STOCK", request.optionId().toString(), "STOCK_RESTORED", event); } + + // 다중 재고 비동기 처리 트랜잭션 + @Transactional + public void processOrderCreation(OrderCreatedMessage msg) { + if (processedEventRepository.existsById(msg.eventId())) { + log.warn("[Idempotency] 이미 처리된 주문 메시지입니다. reservationId={}", msg.reservationId()); + return; + } + + Map groupedItems = msg.items().stream() + .collect(Collectors.groupingBy( + OrderCreatedMessage.OrderItemDto::optionId, + Collectors.summingInt(OrderCreatedMessage.OrderItemDto::quantity) + )); + + List modifiedStocks = new ArrayList<>(); + + // 합산된 수량으로 재고 차감 시도 + for (Map.Entry entry : groupedItems.entrySet()) { + UUID optionId = entry.getKey(); + int quantity = entry.getValue(); + + Stock stock = stockRepository.findById(optionId) + .orElseThrow(StockNotFoundException::new); + + stock.reserve(quantity); + modifiedStocks.add(stock); + + // 카탈로그 동기화용 이벤트 발송 + // 파티션 키는 optionId - 동일 옵션에 대한 차감 순서 FIFO + StockReservedEvent reservedEvent = new StockReservedEvent( + stock.getOptionId(), + stock.getTotalQuantity(), + stock.getCurrentDailyStock() + ); + outboxHelper.append("STOCK", stock.getOptionId().toString(), "STOCK_RESERVED", reservedEvent); + log.info("[Inventory Saga] 다중 주문 옵션 차감 이벤트 적재 완료: optionId={}, 수량={}", stock.getOptionId(), quantity); + + // 품절 처리 이벤트 발송 준비 + if (stock.getTotalQuantity() == 0) { + ProductOption option = productOptionRepository.findById(stock.getOptionId()).orElseThrow(); + Product product = option.getProduct(); + if (product.getStatus() != ProductStatus.SOLDOUT && product.getStatus() != ProductStatus.DELETED + && product.getStatus() != ProductStatus.EXPIRED) { + product.changeStatus(ProductStatus.SOLDOUT); + productRepository.save(product); + outboxHelper.append("PRODUCT", product.getId().toString(), "PRODUCT_STATUS_CHANGED", + new ProductStatusChangedEvent(product.getId(), product.getStatus().name())); + } + } + } + + stockRepository.saveAll(modifiedStocks); + processedEventRepository.save(new ProcessedEvent(msg.eventId())); + + // 모두 성공 시 승인 이벤트 적재 + outboxHelper.append("ORDER", msg.reservationId().toString(), "ORDER_APPROVED", + new OrderApprovedEvent(msg.reservationId())); + log.info("[Inventory Saga] 다중 주문 재고 선점 성공 -> 승인 이벤트 적재 완료: {}", msg.reservationId()); + } } diff --git a/src/main/java/com/michelet/inventory/application/StockLockFacade.java b/src/main/java/com/michelet/inventory/application/StockLockFacade.java index b7f6630..1b5dc2a 100644 --- a/src/main/java/com/michelet/inventory/application/StockLockFacade.java +++ b/src/main/java/com/michelet/inventory/application/StockLockFacade.java @@ -1,7 +1,14 @@ package com.michelet.inventory.application; -import com.michelet.inventory.presentation.dto.ReserveStockRequest; -import com.michelet.inventory.presentation.dto.RestoreStockRequest; +import com.michelet.common.exception.BusinessException; +import com.michelet.inventory.application.dto.ReserveStockRequest; +import com.michelet.inventory.application.dto.RestoreStockRequest; +import com.michelet.inventory.domain.exception.InventoryErrorCode; +import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage; +import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage.OrderItemDto; +import java.util.ArrayList; +import java.util.List; +import java.util.UUID; import java.util.concurrent.TimeUnit; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -67,4 +74,43 @@ public void restoreStockWithLock(RestoreStockRequest request) { } } } + + // 3. 다중 품목 비동기 재고 선점 (OrderCreatedMessage 수신용) + public void reserveOrderStocksWithLock( + OrderCreatedMessage msg) { + // 데드락 방지를 위해 OptionId를 오름차순으로 정렬하여 락을 순차적으로 획득 + List sortedOptionIds = msg.items().stream() + .map(OrderItemDto::optionId) + .sorted() + .toList(); + + List locks = new ArrayList<>(); + try { + for (java.util.UUID id : sortedOptionIds) { + RLock lock = redissonClient.getLock("stock:" + id); + if (lock.tryLock(5, TimeUnit.SECONDS)) { + locks.add(lock); + } else { + log.error("[StockLockFacade] 다중 재고 선점 락 획득 실패 - OptionId: {}", id); + throw new BusinessException( + InventoryErrorCode.CONCURRENCY_ERROR); + } + } + // 모든 락 획득 성공 시 비즈니스 로직 실행 + stockCommandService.processOrderCreation(msg); + + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new BusinessException( + InventoryErrorCode.CONCURRENCY_ERROR); + } finally { + // 데드락 및 리소스 낭비 방지를 위해 역순으로 락 해제 + for (int i = locks.size() - 1; i >= 0; i--) { + RLock lock = locks.get(i); + if (lock != null && lock.isHeldByCurrentThread()) { + lock.unlock(); + } + } + } + } } diff --git a/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java b/src/main/java/com/michelet/inventory/application/dto/ReserveStockRequest.java similarity index 90% rename from src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java rename to src/main/java/com/michelet/inventory/application/dto/ReserveStockRequest.java index 8b337fc..969b629 100644 --- a/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java +++ b/src/main/java/com/michelet/inventory/application/dto/ReserveStockRequest.java @@ -1,4 +1,4 @@ -package com.michelet.inventory.presentation.dto; +package com.michelet.inventory.application.dto; import jakarta.validation.constraints.Min; import jakarta.validation.constraints.NotNull; diff --git a/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java b/src/main/java/com/michelet/inventory/application/dto/RestoreStockRequest.java similarity index 84% rename from src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java rename to src/main/java/com/michelet/inventory/application/dto/RestoreStockRequest.java index acb0540..1fbd997 100644 --- a/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java +++ b/src/main/java/com/michelet/inventory/application/dto/RestoreStockRequest.java @@ -1,4 +1,4 @@ -package com.michelet.inventory.presentation.dto; +package com.michelet.inventory.application.dto; import jakarta.validation.constraints.NotNull; import jakarta.validation.constraints.Positive; diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java index 581c40f..fd54db2 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java @@ -19,7 +19,8 @@ public class DeadLetterConsumer { */ @KafkaListener( topics = { - "${inventory.kafka.topic.restore-request:order.stock-restore.requested}.DLT", + // 개별로 나뉘어 있던 토픽을 지우고, 통합 커맨드 토픽의 DLT 구독 + "${inventory.kafka.topic.inventory-command:inventory.command}.DLT", "${inventory.kafka.topic.restored:stock.restored}.DLT", "${inventory.kafka.topic.reserved:stock.reserved}.DLT" }, 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 9f9563f..eff66b6 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java @@ -1,24 +1,71 @@ package com.michelet.inventory.infrastructure.messaging; +import com.michelet.common.exception.BusinessException; +import com.michelet.inventory.application.InventoryOutboxHelper; import com.michelet.inventory.application.StockLockFacade; +import com.michelet.inventory.application.dto.RestoreStockRequest; +import com.michelet.inventory.domain.exception.ConcurrencyFailureException; +import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage; +import com.michelet.inventory.infrastructure.messaging.dto.OrderRejectedEvent; import com.michelet.inventory.infrastructure.messaging.dto.StockRestoreMessage; -import com.michelet.inventory.presentation.dto.RestoreStockRequest; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.kafka.annotation.KafkaHandler; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; @Slf4j @Component @RequiredArgsConstructor +// 클래스 레벨에서 단일 토픽(inventory.command)을 구독하도록 통합 +@KafkaListener( + topics = "${inventory.kafka.topic.inventory-command:inventory.command}", + groupId = "${spring.kafka.consumer.group-id:inventory-service-consumer}" +) public class OrderEventConsumer { private final StockLockFacade stockLockFacade; + private final InventoryOutboxHelper outboxHelper; - @KafkaListener( - topics = "${inventory.kafka.topic.restore-request:order.stock-restore.requested}", - groupId = "${spring.kafka.consumer.group-id:inventory-service-consumer}" - ) + // TypeId 헤더를 보고 스프링 카프카가 이 메서드로 라우팅해 줌 + @KafkaHandler + public void consumeOrderCreated(OrderCreatedMessage payload) { + // payload null 체크 (Tombstone 메시지 방어 - NPE 방지) + if (payload == null) { + log.error("[Kafka Consumer] 잘못된 주문 생성 이벤트 수신: payload가 null입니다 (Tombstone 메시지일 가능성)."); + throw new IllegalArgumentException("주문 생성 이벤트 파싱 오류: payload는 null일 수 없습니다."); + } + + // null 체크 이후에 로깅 (In-Order 처리) + log.info("[Kafka Consumer] 신규 주문 생성 메시지 수신 -> 재고 다중 차감 시도 (In-Order 처리): reservationId={}", + payload.reservationId()); + try { + stockLockFacade.reserveOrderStocksWithLock(payload); + } catch (IllegalArgumentException e) { + log.error("[Kafka Consumer] 비즈니스/검증 룰 위반 에러 (DLT 직행 대상). reservationId: {}", payload.reservationId(), e); + throw e; + + } catch (ConcurrencyFailureException e) { + // 여기서 동시성 에러를 가로채서 재시도 처리함! + // 락 획득 실패(일시적 경합)인 경우, 주문 거절을 하지 않고 RuntimeException을 던져 카프카 재시도를 유도 + log.error("[Kafka Consumer] 락 경합으로 인한 일시적 실패 (재시도 대상). reservationId={}", payload.reservationId(), e); + throw new RuntimeException("재고 락 경합으로 인한 주문 처리 지연", e); + + } catch (BusinessException e) { + // 품절, 한도 초과 등 명확한 비즈니스 예외 시에만 오더 서비스에 거절 이벤트 전송 + log.warn("[Kafka Consumer] 비즈니스 로직에 의한 재고 차감 실패 (거절 이벤트 정상 발행). 사유: {}", e.getMessage()); + // 롤백된 트랜잭션 밖에서 아웃박스를 안전하게 저장하기 위해 appendIndependent 사용! + outboxHelper.appendIndependent("ORDER", payload.reservationId().toString(), "ORDER_REJECTED", + new OrderRejectedEvent(payload.reservationId(), e.getMessage())); + + } catch (Exception e) { + log.error("[Kafka Consumer] 재고 차감 중 일시적 에러 발생 (재시도 대상). reservationId: {}", payload.reservationId(), e); + throw new RuntimeException("재고 다중 차감 실패", e); + } + } + + // 복구 메시지가 들어오면 이 메서드로 라우팅해 줌 + @KafkaHandler public void consumeStockRestoreRequest(StockRestoreMessage payload) { if (payload == null) { log.error("[Kafka Consumer] 잘못된 복구 이벤트 수신: payload가 null입니다 (Tombstone 메시지일 가능성)."); @@ -26,7 +73,7 @@ public void consumeStockRestoreRequest(StockRestoreMessage payload) { throw new IllegalArgumentException("재고 복구 이벤트 파싱 오류: payload는 null일 수 없습니다."); } - log.info("[Kafka Consumer] 재고 복구 이벤트 수신: eventId={}, optionId={}, quantity={}", payload.eventId(), + log.info("[Kafka Consumer] 재고 복구 이벤트 수신 (In-Order 처리): eventId={}, optionId={}, quantity={}", payload.eventId(), payload.optionId(), payload.quantity()); try { @@ -49,4 +96,14 @@ public void consumeStockRestoreRequest(StockRestoreMessage payload) { throw new RuntimeException("재고 복구 컨슈머 처리 실패", e); } } + + // 알 수 없는 타입의 객체가 inventory.command로 들어올 경우, 조용히 넘기지 않고 예외를 발생시켜 DLT로 격리 + @KafkaHandler(isDefault = true) + public void unknown(Object object) { + log.error("[Kafka Consumer] 지원하지 않는 커맨드 타입 수신. object={}", object); + throw new IllegalArgumentException( + "지원하지 않는 inventory.command payload 타입: " + + (object == null ? "null" : object.getClass().getName()) + ); + } } diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderApprovedEvent.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderApprovedEvent.java new file mode 100644 index 0000000..f2a524c --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderApprovedEvent.java @@ -0,0 +1,11 @@ +package com.michelet.inventory.infrastructure.messaging.dto; + +import java.util.UUID; + +public record OrderApprovedEvent(UUID reservationId) { + public OrderApprovedEvent { + if (reservationId == null) { + throw new IllegalArgumentException("reservationId는 null일 수 없습니다."); + } + } +} diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderCreatedMessage.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderCreatedMessage.java new file mode 100644 index 0000000..f7ea7b2 --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderCreatedMessage.java @@ -0,0 +1,40 @@ +package com.michelet.inventory.infrastructure.messaging.dto; + +import java.util.List; +import java.util.UUID; + +public record OrderCreatedMessage( + UUID eventId, + UUID reservationId, + List items +) { + public OrderCreatedMessage { + if (eventId == null) { + throw new IllegalArgumentException("메시지 파싱 오류: eventId는 null일 수 없습니다."); + } + if (reservationId == null) { + throw new IllegalArgumentException("메시지 파싱 오류: reservationId는 null일 수 없습니다."); + } + if (items == null || items.isEmpty()) { + throw new IllegalArgumentException("메시지 파싱 오류: 주문 항목은 1개 이상이어야 합니다."); + } + // 리스트 내부 null 원소 검증 + if (items.stream().anyMatch(java.util.Objects::isNull)) { + throw new IllegalArgumentException("메시지 파싱 오류: 주문 항목에 null 값이 포함될 수 없습니다."); + } + + // 외부 조작 방지를 위한 불변 리스트 복사 + items = List.copyOf(items); + } + + public record OrderItemDto(UUID optionId, Integer quantity) { + public OrderItemDto { + if (optionId == null) { + throw new IllegalArgumentException("메시지 파싱 오류: optionId는 null일 수 없습니다."); + } + if (quantity == null || quantity <= 0) { + throw new IllegalArgumentException("메시지 파싱 오류: 수량은 1 이상이어야 합니다."); + } + } + } +} diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderRejectedEvent.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderRejectedEvent.java new file mode 100644 index 0000000..115e0ea --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderRejectedEvent.java @@ -0,0 +1,14 @@ +package com.michelet.inventory.infrastructure.messaging.dto; + +import java.util.UUID; + +public record OrderRejectedEvent(UUID reservationId, String reason) { + public OrderRejectedEvent { + if (reservationId == null) { + throw new IllegalArgumentException("reservationId는 null일 수 없습니다."); + } + if (reason == null || reason.isBlank()) { + throw new IllegalArgumentException("거절 사유(reason)는 필수입니다."); + } + } +} diff --git a/src/main/java/com/michelet/inventory/presentation/InternalStockController.java b/src/main/java/com/michelet/inventory/presentation/InternalStockController.java deleted file mode 100644 index 030b85d..0000000 --- a/src/main/java/com/michelet/inventory/presentation/InternalStockController.java +++ /dev/null @@ -1,41 +0,0 @@ -package com.michelet.inventory.presentation; - -import com.michelet.common.response.ApiResponse; -import com.michelet.inventory.application.StockLockFacade; -import com.michelet.inventory.presentation.dto.ReserveStockRequest; -import com.michelet.inventory.presentation.dto.RestoreStockRequest; -import jakarta.validation.Valid; -import lombok.RequiredArgsConstructor; -import org.springframework.http.ResponseEntity; -import org.springframework.web.bind.annotation.PostMapping; -import org.springframework.web.bind.annotation.RequestBody; -import org.springframework.web.bind.annotation.RequestMapping; -import org.springframework.web.bind.annotation.RestController; - -@RestController -@RequestMapping("/internal/stocks") -@RequiredArgsConstructor -public class InternalStockController { - - private final StockLockFacade stockLockFacade; - - @PostMapping("/reserve") - public ResponseEntity> reserveStock( - @RequestBody @Valid ReserveStockRequest request - ) { - stockLockFacade.reserveStockWithLock(request); - 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 - ) { - stockLockFacade.restoreStockWithLock(request); - return ResponseEntity.ok(ApiResponse.ok(null)); - } -} diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index aebc366..ef071bb 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -32,11 +32,10 @@ spring: properties: # 실제 역직렬화를 수행할 델리게이트 클래스 지정 spring.deserializer.value.delegate.class: org.springframework.kafka.support.serializer.JsonDeserializer - spring.json.trusted.packages: "com.michelet.inventory.infrastructure.messaging.dto,com.michelet.order.application.dto" + spring.json.trusted.packages: "com.michelet.inventory.infrastructure.messaging.dto" spring.json.use.type.headers: true # 서로 다른 패키지의 이벤트를 1:1로 매핑 - spring.json.type.mapping: "com.michelet.order.application.dto.StockRestoreEventPayload:com.michelet.inventory.infrastructure.messaging.dto.StockRestoreMessage" - + spring.json.type.mapping: "com.michelet.order.application.dto.StockRestoreEventPayload:com.michelet.inventory.infrastructure.messaging.dto.StockRestoreMessage,com.michelet.order.application.dto.OrderCreatedEventPayload:com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage" server: port: 19900 @@ -48,11 +47,14 @@ inventory: topic: product-created: "product.created" reserved: "stock.reserved" - restore-request: "order.stock-restore.requested" # 오더의 '명령'을 수신할 토픽 restored: "stock.restored" # 카탈로그로 '결과'를 발송할 토픽 status-changed: "product.status-changed" daily-reset: "stock.daily-reset" product-updated: "product.updated" + # 분리되어 있던 토픽을 지우고, 단일 통합 커맨드 토픽 추가 + inventory-command: "inventory.command" + order-approved: "order.approved" + order-rejected: "order.rejected" consumer: retry: interval-ms: 1000 diff --git a/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java b/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java index b9daac7..1e31839 100644 --- a/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java +++ b/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java @@ -9,6 +9,8 @@ import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; +import com.michelet.inventory.application.dto.ReserveStockRequest; +import com.michelet.inventory.application.dto.RestoreStockRequest; import com.michelet.inventory.application.dto.StockReservedEvent; import com.michelet.inventory.application.dto.StockRestoredEvent; import com.michelet.inventory.domain.exception.MaxLimitExceededException; @@ -20,8 +22,6 @@ import com.michelet.inventory.domain.repository.ProductOptionRepository; import com.michelet.inventory.domain.repository.ProductRepository; import com.michelet.inventory.domain.repository.StockRepository; -import com.michelet.inventory.presentation.dto.ReserveStockRequest; -import com.michelet.inventory.presentation.dto.RestoreStockRequest; import java.util.Optional; import java.util.UUID; import org.junit.jupiter.api.DisplayName; diff --git a/src/test/java/com/michelet/inventory/application/StockConcurrencyIntegrationTest.java b/src/test/java/com/michelet/inventory/application/StockConcurrencyIntegrationTest.java index d6b4ee4..ba7815d 100644 --- a/src/test/java/com/michelet/inventory/application/StockConcurrencyIntegrationTest.java +++ b/src/test/java/com/michelet/inventory/application/StockConcurrencyIntegrationTest.java @@ -5,10 +5,10 @@ import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.BDDMockito.given; +import com.michelet.inventory.application.dto.ReserveStockRequest; import com.michelet.inventory.domain.model.Stock; import com.michelet.inventory.domain.repository.StockRepository; import com.michelet.inventory.infrastructure.repository.JpaStockRepository; -import com.michelet.inventory.presentation.dto.ReserveStockRequest; import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; diff --git a/src/test/resources/application-test.yml b/src/test/resources/application-test.yml index c334cb0..1376627 100644 --- a/src/test/resources/application-test.yml +++ b/src/test/resources/application-test.yml @@ -51,4 +51,4 @@ inventory: restored: "stock.restored" status-changed: "product.status-changed" daily-reset: "stock.daily-reset" - + inventory-command: "inventory.command"