diff --git a/docker-compose.yml b/docker-compose.yml index bde06f1..7332929 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -23,6 +23,7 @@ services: RABBITMQ_PREFETCH: ${RABBITMQ_PREFETCH:-250} MOCK_PROVIDER_FORCE_FAILURE: ${MOCK_PROVIDER_FORCE_FAILURE:-false} MOCK_PROVIDER_FAIL_FIRST_ATTEMPTS: ${MOCK_PROVIDER_FAIL_FIRST_ATTEMPTS:-0} + OUTBOX_PUBLISH_INTERVAL_MS: ${OUTBOX_PUBLISH_INTERVAL_MS:-1000} healthcheck: test: diff --git a/src/main/java/com/backendsystemdesignlab/notification/NotificationSystemApplication.java b/src/main/java/com/backendsystemdesignlab/notification/NotificationSystemApplication.java index 21332b5..981e162 100644 --- a/src/main/java/com/backendsystemdesignlab/notification/NotificationSystemApplication.java +++ b/src/main/java/com/backendsystemdesignlab/notification/NotificationSystemApplication.java @@ -2,8 +2,10 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.scheduling.annotation.EnableScheduling; @SpringBootApplication +@EnableScheduling public class NotificationSystemApplication { public static void main(String[] args) { diff --git a/src/main/java/com/backendsystemdesignlab/notification/messaging/DeliveryMessageHandler.java b/src/main/java/com/backendsystemdesignlab/notification/messaging/DeliveryMessageHandler.java index 9ec3814..8d9bdc9 100644 --- a/src/main/java/com/backendsystemdesignlab/notification/messaging/DeliveryMessageHandler.java +++ b/src/main/java/com/backendsystemdesignlab/notification/messaging/DeliveryMessageHandler.java @@ -23,6 +23,8 @@ public class DeliveryMessageHandler { public void handle(DeliveryMessage message) { + if (transactionService.isAlreadyProcessed(message.deliveryId())) return; + boolean success = send(message); if (success) { @@ -43,12 +45,14 @@ public void handle(DeliveryMessage message) { private boolean send(DeliveryMessage message) { + String idempotencyKey = "notification-delivery-" + message.deliveryId(); + try { ProviderResult result = switch (message.channel()) { - case PUSH -> pushProvider.send(message.destination()); - case SMS -> smsProvider.send(message.destination()); - case EMAIL -> emailProvider.send(message.destination()); + case PUSH -> pushProvider.send(message.destination(), idempotencyKey); + case SMS -> smsProvider.send(message.destination(), idempotencyKey); + case EMAIL -> emailProvider.send(message.destination(), idempotencyKey); }; return result.success(); } catch (RuntimeException e) { diff --git a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/EmailProvider.java b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/EmailProvider.java index d57d1a7..1290e5d 100644 --- a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/EmailProvider.java +++ b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/EmailProvider.java @@ -1,5 +1,5 @@ package com.backendsystemdesignlab.notification.notification.provider; public interface EmailProvider { - ProviderResult send(String email); + ProviderResult send(String email, String idempotencyKey); } diff --git a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/PushProvider.java b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/PushProvider.java index 78c7aa3..9bb4dd1 100644 --- a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/PushProvider.java +++ b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/PushProvider.java @@ -1,5 +1,5 @@ package com.backendsystemdesignlab.notification.notification.provider; public interface PushProvider { - ProviderResult send(String deviceToken); + ProviderResult send(String deviceToken, String idempotencyKey); } diff --git a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/SmsProvider.java b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/SmsProvider.java index 6aac852..037cbbb 100644 --- a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/SmsProvider.java +++ b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/SmsProvider.java @@ -1,5 +1,5 @@ package com.backendsystemdesignlab.notification.notification.provider; public interface SmsProvider { - ProviderResult send(String phoneNumber); + ProviderResult send(String phoneNumber, String idempotencyKey); } diff --git a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockEmailProvider.java b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockEmailProvider.java index 70b4403..b399558 100644 --- a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockEmailProvider.java +++ b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockEmailProvider.java @@ -6,6 +6,7 @@ import org.springframework.stereotype.Component; import java.util.Map; +import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; @@ -14,6 +15,7 @@ public class MockEmailProvider implements EmailProvider { private final boolean forceFailure; private final int failFirstAttempts; + private final Set processedKeys = ConcurrentHashMap.newKeySet(); private final Map attempts = new ConcurrentHashMap<>(); @@ -25,8 +27,13 @@ public MockEmailProvider( } @Override - public ProviderResult send(String email) { + public ProviderResult send(String email, String idempotencyKey) { simulateDelay(); + + if (processedKeys.contains(idempotencyKey)) { + return new ProviderResult(true); + } + if (forceFailure) return new ProviderResult(false); int currentAttempt = attempts.computeIfAbsent(email, key -> new AtomicInteger()).incrementAndGet(); // 증가 연산을 원자적 처리 @@ -35,6 +42,7 @@ public ProviderResult send(String email) { return new ProviderResult(false); } + processedKeys.add(idempotencyKey); return new ProviderResult(true); } diff --git a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockPushProvider.java b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockPushProvider.java index 1ac2cc7..92a136e 100644 --- a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockPushProvider.java +++ b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockPushProvider.java @@ -5,19 +5,29 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; + @Component public class MockPushProvider implements PushProvider { private final boolean forceFailure; + private final Set processedKeys = ConcurrentHashMap.newKeySet(); public MockPushProvider(@Value("${mock.provider.force-failure:false}") boolean forceFailure) { this.forceFailure = forceFailure; } @Override - public ProviderResult send(String deviceToken) { + public ProviderResult send(String deviceToken, String idempotencyKey) { simulateDelay(); + + if (processedKeys.contains(idempotencyKey)) { + return new ProviderResult(true); + } + if (forceFailure) return new ProviderResult(false); + processedKeys.add(idempotencyKey); return new ProviderResult(true); } diff --git a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockSmsProvider.java b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockSmsProvider.java index ef1681f..98212c4 100644 --- a/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockSmsProvider.java +++ b/src/main/java/com/backendsystemdesignlab/notification/notification/provider/mock/MockSmsProvider.java @@ -5,19 +5,29 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; + @Component public class MockSmsProvider implements SmsProvider { private final boolean forceFailure; + private final Set processedKeys = + ConcurrentHashMap.newKeySet(); public MockSmsProvider(@Value("${mock.provider.force-failure:false}") boolean forceFailure) { this.forceFailure = forceFailure; } @Override - public ProviderResult send(String deviceToken) { + public ProviderResult send(String phoneNumber, String idempotencyKey) { simulateDelay(); + if (processedKeys.contains(idempotencyKey)) { + return new ProviderResult(true); + } if (forceFailure) return new ProviderResult(false); + + processedKeys.add(idempotencyKey); return new ProviderResult(true); } diff --git a/src/main/java/com/backendsystemdesignlab/notification/notification/service/NotificationService.java b/src/main/java/com/backendsystemdesignlab/notification/notification/service/NotificationService.java index eaeea0f..00e0607 100644 --- a/src/main/java/com/backendsystemdesignlab/notification/notification/service/NotificationService.java +++ b/src/main/java/com/backendsystemdesignlab/notification/notification/service/NotificationService.java @@ -28,7 +28,7 @@ public SendNotificationResponse send(SendNotificationRequest request) { ); } - deliveryPublisher.publishAll(prepared.notificationId(), prepared.deliveries()); +// deliveryPublisher.publishAll(prepared.notificationId(), prepared.deliveries()); return new SendNotificationResponse( prepared.notificationId(), diff --git a/src/main/java/com/backendsystemdesignlab/notification/notification/service/NotificationTransactionService.java b/src/main/java/com/backendsystemdesignlab/notification/notification/service/NotificationTransactionService.java index 363224c..9921b65 100644 --- a/src/main/java/com/backendsystemdesignlab/notification/notification/service/NotificationTransactionService.java +++ b/src/main/java/com/backendsystemdesignlab/notification/notification/service/NotificationTransactionService.java @@ -8,6 +8,8 @@ import com.backendsystemdesignlab.notification.notification.dto.SendNotificationRequest; import com.backendsystemdesignlab.notification.notification.repository.NotificationDeliveryRepository; import com.backendsystemdesignlab.notification.notification.repository.NotificationRepository; +import com.backendsystemdesignlab.notification.outbox.OutboxEvent; +import com.backendsystemdesignlab.notification.outbox.OutboxEventRepository; import com.backendsystemdesignlab.notification.user.domain.NotificationChannel; import com.backendsystemdesignlab.notification.user.domain.NotificationPreference; import com.backendsystemdesignlab.notification.user.domain.User; @@ -22,6 +24,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Set; +import java.util.UUID; import java.util.stream.Collectors; @Service @@ -33,6 +36,7 @@ public class NotificationTransactionService { private final NotificationPreferenceRepository preferenceRepository; private final NotificationRepository notificationRepository; private final NotificationDeliveryRepository deliveryRepository; + private final OutboxEventRepository outboxEventRepository; @Transactional public PreparedNotification prepare(SendNotificationRequest request) { @@ -93,6 +97,19 @@ public PreparedNotification prepare(SendNotificationRequest request) { deliveryRepository.flush(); // DB의 ID를 얻기 위함 (delivery.getId()) + for (NotificationDelivery delivery : deliveries) { + + OutboxEvent outboxEvent = new OutboxEvent( + UUID.randomUUID().toString(), + notification.getId(), + delivery.getId(), + delivery.getChannel(), + delivery.getDestination() + ); + + outboxEventRepository.save(outboxEvent); + } + notification.startProcessing(); List commands = deliveries.stream() @@ -216,4 +233,29 @@ private void createEmailDelivery(User user, Notification notification, List new IllegalArgumentException("Delivery not found: " + deliveryId)); + + return delivery.getStatus() == DeliveryStatus.SENT || delivery.getStatus() == DeliveryStatus.FAILED; + } + + @Transactional + public void recordPublishFinalFailure(Long notificationId, Long deliveryId) { + Notification notification = notificationRepository.findByIdForUpdate(notificationId) + .orElseThrow(() -> new IllegalArgumentException("알림을 찾을 수 없습니다.")); + + NotificationDelivery delivery = deliveryRepository.findById(deliveryId) + .orElseThrow(() -> new IllegalArgumentException("전송 정보를 찾을 수 없습니다.")); + + if (delivery.getStatus() != DeliveryStatus.PENDING) { + return; + } + + delivery.markFailed(); + + updateNotificationStatus(notificationId, notification); + } } diff --git a/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxEvent.java b/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxEvent.java new file mode 100644 index 0000000..7f091ce --- /dev/null +++ b/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxEvent.java @@ -0,0 +1,89 @@ +package com.backendsystemdesignlab.notification.outbox; + +import com.backendsystemdesignlab.notification.user.domain.NotificationChannel; +import jakarta.persistence.*; +import lombok.AccessLevel; +import lombok.Getter; +import lombok.NoArgsConstructor; + +import java.time.LocalDateTime; + +@Entity +@Table( + name = "outbox_events", + indexes = { + @Index( + name = "idx_outbox_status_created_at", + columnList = "status, created_at" + ) + } +) +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +public class OutboxEvent { + + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @Column(nullable = false, unique = true) + private String messageId; + + @Column(nullable = false) + private Long notificationId; + + @Column(nullable = false) + private Long deliveryId; + + @Enumerated(EnumType.STRING) + @Column(nullable = false) + private NotificationChannel channel; + + @Column(nullable = false) + private String destination; + + @Enumerated(EnumType.STRING) + @Column(nullable = false) + private OutboxStatus status; + + @Column(nullable = false) + private LocalDateTime createdAt; + + private LocalDateTime publishedAt; + + @Column(nullable = false) + private int attemptCount; + + private String lastError; + + public OutboxEvent( + String messageId, + Long notificationId, + Long deliveryId, + NotificationChannel channel, + String destination + ) { + this.messageId = messageId; + this.notificationId = notificationId; + this.deliveryId = deliveryId; + this.channel = channel; + this.destination = destination; + this.status = OutboxStatus.PENDING; + this.attemptCount = 0; + this.createdAt = LocalDateTime.now(); + } + + public void markPublished() { + this.status = OutboxStatus.PUBLISHED; + this.publishedAt = LocalDateTime.now(); + } + + public void recordFailure(String error) { + this.attemptCount++; + this.lastError = error; + } + + public void markFailed() { + this.status = OutboxStatus.FAILED; + } +} diff --git a/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxEventRepository.java b/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxEventRepository.java new file mode 100644 index 0000000..dc493d3 --- /dev/null +++ b/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxEventRepository.java @@ -0,0 +1,10 @@ +package com.backendsystemdesignlab.notification.outbox; + +import org.springframework.data.jpa.repository.JpaRepository; + +import java.util.List; + +public interface OutboxEventRepository extends JpaRepository { + + List findTop100ByStatusOrderByCreatedAtAsc(OutboxStatus status); +} diff --git a/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxPublisher.java b/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxPublisher.java new file mode 100644 index 0000000..20bbe4f --- /dev/null +++ b/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxPublisher.java @@ -0,0 +1,111 @@ +package com.backendsystemdesignlab.notification.outbox; + +import com.backendsystemdesignlab.notification.messaging.DeliveryMessage; +import com.backendsystemdesignlab.notification.messaging.RabbitMqConfig; +import com.backendsystemdesignlab.notification.notification.service.NotificationTransactionService; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.amqp.rabbit.connection.CorrelationData; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +import java.util.List; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; + +@Slf4j +@Component +@RequiredArgsConstructor +public class OutboxPublisher { + + private final OutboxTransactionService transactionService; + private final RabbitTemplate rabbitTemplate; + private final NotificationTransactionService notificationTransactionService; + + // DB Connection 없음 (RabbitMQ가 DB를 잡고 있지 않게) + @Scheduled( + fixedDelayString = "${outbox.publish-interval-ms:1000}" + ) + public void publishPendingEvents() { + List events = transactionService.findPendingEvents(); + + for (OutboxEvent event : events) { + publish(event); + } + } + + private void publish(OutboxEvent event) { + try { + DeliveryMessage message = new DeliveryMessage( + event.getNotificationId(), + event.getDeliveryId(), + event.getChannel(), + event.getDestination(), + 1 + ); + + CorrelationData correlationData = new CorrelationData(String.valueOf(event.getId())); + + rabbitTemplate.convertAndSend( + RabbitMqConfig.EXCHANGE, + routingKey(event), + message, + correlationData + ); + + // Spring AMQP < CorrelationData < CompletableFuture (Confirm의 결과) + CorrelationData.Confirm confirm = correlationData.getFuture().get(5, TimeUnit.SECONDS); // 5초 기다림 + + if (!confirm.ack()) { // RabbitMQ Broker가 잘 받았는지 (Publisher -> Broker) + handleFailure(event, "NACK: " + confirm.reason()); + return; + } + + if (correlationData.getReturned() != null) { // Broker Exchange -> Queue로 라우팅됐는가? (라우팅 실패) 예: 잘못된 라우팅키 + String error = "RETURN: " + correlationData.getReturned().getReplyText(); + handleFailure(event, error); + return; + } + + + // @Transaction + transactionService.markPublished(event.getId()); + + log.debug("Outbox published. outboxId={}, deliveryId={}", event.getId(), event.getDeliveryId()); + } catch (TimeoutException e) { + handleFailure(event, "CONFIRM_TIMEOUT"); + } catch (Exception e) { + handleFailure(event, e.getMessage()); + } + } + + private String routingKey(OutboxEvent event) { + return switch (event.getChannel()) { + case PUSH -> + RabbitMqConfig.PUSH_ROUTING_KEY; + + case SMS -> + RabbitMqConfig.SMS_ROUTING_KEY; + + case EMAIL -> + RabbitMqConfig.EMAIL_ROUTING_KEY; + }; + } + + private void handleFailure(OutboxEvent event, String reason) { + + boolean finalFailure = transactionService.recordPublishFailure(event.getId(), reason); + + if (finalFailure) { + notificationTransactionService.recordPublishFinalFailure( + event.getNotificationId(), + event.getDeliveryId() + ); + + log.error("Outbox 마지막 시도 실패. outboxId={}, reason={}", event.getId(), reason); + } else { + log.warn("Outbox publish 실패. outboxId={}, reason={}", event.getId(), reason); + } + } +} diff --git a/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxStatus.java b/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxStatus.java new file mode 100644 index 0000000..4f2f6be --- /dev/null +++ b/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxStatus.java @@ -0,0 +1,7 @@ +package com.backendsystemdesignlab.notification.outbox; + +public enum OutboxStatus { + PENDING, + PUBLISHED, + FAILED +} diff --git a/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxTransactionService.java b/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxTransactionService.java new file mode 100644 index 0000000..5ce7d6a --- /dev/null +++ b/src/main/java/com/backendsystemdesignlab/notification/outbox/OutboxTransactionService.java @@ -0,0 +1,50 @@ +package com.backendsystemdesignlab.notification.outbox; + +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +import java.util.List; + +@Service +@RequiredArgsConstructor +public class OutboxTransactionService { + + private final OutboxEventRepository outboxEventRepository; + private static final int MAX_PUBLISH_ATTEMPTS = 5; + + @Transactional(readOnly = true) + public List findPendingEvents() { + return outboxEventRepository.findTop100ByStatusOrderByCreatedAtAsc(OutboxStatus.PENDING); + } + + @Transactional + public void markPublished(Long outboxEventId) { + OutboxEvent event = outboxEventRepository.findById(outboxEventId) + .orElseThrow(() -> new IllegalArgumentException("Outbox event not found: " + outboxEventId)); + + if (event.getStatus() == OutboxStatus.PUBLISHED) { + return; + } + event.markPublished(); + } + + @Transactional + public boolean recordPublishFailure(Long outboxEventId, String error) { + OutboxEvent event = outboxEventRepository.findById(outboxEventId) + .orElseThrow(() -> new IllegalArgumentException("Outbox event not found: " + outboxEventId)); + + if (event.getStatus() != OutboxStatus.PENDING) { + return false; + } + + event.recordFailure(error); + + if (event.getAttemptCount() >= MAX_PUBLISH_ATTEMPTS) { + event.markFailed(); + return true; + } + + return false; + } +} diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index 1e27963..d74be54 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -27,6 +27,10 @@ spring: listener: simple: prefetch: ${RABBITMQ_PREFETCH:250} + publisher-confirm-type: correlated + publisher-returns: true + template: + mandatory: true management: endpoints: @@ -53,4 +57,7 @@ management: mock: provider: force-failure: ${MOCK_PROVIDER_FORCE_FAILURE:false} - fail-first-attempts: ${MOCK_PROVIDER_FAIL_FIRST_ATTEMPTS:0} \ No newline at end of file + fail-first-attempts: ${MOCK_PROVIDER_FAIL_FIRST_ATTEMPTS:0} + +outbox: + publish-interval-ms: ${OUTBOX_PUBLISH_INTERVAL_MS:1000} \ No newline at end of file