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
@@ -0,0 +1,124 @@
package com.momogo.api.room.consumer;

import com.momogo.core.common.config.KafkaTopics;
import com.momogo.core.common.exception.BusinessException;
import com.momogo.core.domain.room.exception.RoomErrorCode;
import com.momogo.core.domain.room.entity.RoomProblem;
import com.momogo.core.domain.room.entity.RoomUser;
import com.momogo.core.domain.room.entity.RoomUserId;
import com.momogo.core.domain.room.entity.UserRoomAnswer;
import com.momogo.core.domain.room.event.RoomSubmitEventMessage;
import com.momogo.core.domain.room.repository.RoomProblemRepository;
import com.momogo.core.domain.room.repository.RoomUserRepository;
import com.momogo.core.domain.room.repository.UserRoomAnswerRepository;
import com.momogo.core.domain.user.entity.User;
import jakarta.persistence.EntityManager;
import java.time.Duration;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.function.Function;
import java.util.stream.Collectors;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.dao.DataIntegrityViolationException;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.transaction.support.TransactionSynchronization;
import org.springframework.transaction.support.TransactionSynchronizationManager;

/*
* RoomServiceImpl이 발행한 답안 제출 메시지를 수신하여 실제 DB 저장을 비동기로 수행
* Redis 멱등성 체크를 통해 중복 전달(At-least-once)로 인한 UQ 충돌 방지
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class RoomSubmitKafkaConsumer {

private static final String DEDUP_KEY_PREFIX = "room:submit:processed:";
private static final Duration DEDUP_TTL = Duration.ofHours(24);

private final UserRoomAnswerRepository userRoomAnswerRepository;
private final RoomUserRepository roomUserRepository;
private final RoomProblemRepository roomProblemRepository;
private final RedisTemplate<String, Object> redisTemplate;
private final EntityManager entityManager;

@KafkaListener(topics = KafkaTopics.ROOM_SUBMIT_EVENTS)
@Transactional
public void consume(RoomSubmitEventMessage message) {
log.info("[RoomSubmitKafkaConsumer] 카프카 답안 수신 - eventId: {}, userId: {}, roomId: {}",
message.eventId(), message.userId(), message.roomId());

// 1. Redis atomic setIfAbsent로 "PROCESSING" 상태 선점 (원자적 멱등성 검사)
String dedupKey = DEDUP_KEY_PREFIX + message.eventId();
Boolean isAcquired = redisTemplate.opsForValue().setIfAbsent(dedupKey, "PROCESSING", Duration.ofMinutes(5));
if (Boolean.FALSE.equals(isAcquired)) {
log.warn("[RoomSubmitKafkaConsumer] 이미 처리 중이거나 완료된 답안 제출 이벤트, 중복 스킵 - eventId: {}", message.eventId());
return;
}

// 2. DB 트랜잭션 종결 상태별 후처리 (커밋 성공 시 DONE 갱신, 예외/롤백 시 키 삭제로 카프카 재시도 보장)
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCommit() {
redisTemplate.opsForValue().set(dedupKey, "DONE", DEDUP_TTL);
}

@Override
public void afterCompletion(int status) {
if (status != STATUS_COMMITTED) {
redisTemplate.delete(dedupKey);
redisTemplate.delete("room:submit:claimed:" + message.roomId() + ":" + message.userId());
}
}
});

// 3. 추가 SELECT 없이 EntityManager 프록시 참조 생성
User userProxy = entityManager.getReference(User.class, message.userId());

// 4. 문제 목록 조회 및 해당 방 소속 검증
List<UUID> problemIds = message.answers().stream()
.map(ans -> ans.roomProblemId())
.toList();

Map<UUID, RoomProblem> problemsById = roomProblemRepository.findAllById(problemIds).stream()
.filter(p -> p.getRoom().getId().equals(message.roomId()))
.collect(Collectors.toMap(RoomProblem::getId, Function.identity()));

if (problemsById.size() != problemIds.size()) {
log.error("[RoomSubmitKafkaConsumer] 유효하지 않거나 해당 방에 속하지 않는 문제 ID 포함, 처리 중단 - eventId: {}, roomId: {}, 요청 문제 수: {}, 유효 문제 수: {}",
message.eventId(), message.roomId(), problemIds.size(), problemsById.size());
throw new BusinessException(RoomErrorCode.PROBLEM_NOT_FOUND);
}

List<UserRoomAnswer> userRoomAnswers = message.answers().stream()
.map(ans -> UserRoomAnswer.of(userProxy, problemsById.get(ans.roomProblemId()), ans.userAnswer(), null))
.toList();

try {
// 5. DB 일괄 저장
userRoomAnswerRepository.saveAll(userRoomAnswers);

// 6. 응시 완료 상태 업데이트
RoomUserId roomUserId = new RoomUserId(message.roomId(), message.userId());
roomUserRepository.findById(roomUserId)
.orElseThrow(() -> {
log.error("[RoomSubmitKafkaConsumer] 참여자 정보 없음, 응시 상태 갱신 불가 - userId: {}, roomId: {}",
message.userId(), message.roomId());
return new BusinessException(RoomErrorCode.NOT_ROOM_PARTICIPANT);
})
.attend();
} catch (DataIntegrityViolationException e) {
log.warn("[RoomSubmitKafkaConsumer] DB 답안 중복 유니크 제약조건 충돌, 처리를 정상 스킵함 - eventId: {}, userId: {}, roomId: {}",
message.eventId(), message.userId(), message.roomId());
return;
}

log.info("[RoomSubmitKafkaConsumer] DB 저장 완료 - userId: {}, roomId: {}, 문항 수: {}",
message.userId(), message.roomId(), userRoomAnswers.size());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ public final class KafkaTopics {

public static final String NOTIFICATION_EVENTS = "notification-events";
public static final String AI_GRADING_EVENTS = "ai-grading-events";
public static final String ROOM_SUBMIT_EVENTS = "room-submit-events";

private KafkaTopics() {
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
package com.momogo.core.domain.room.event;

import com.momogo.core.domain.room.dto.request.ProblemAnswerRequest;
import java.time.OffsetDateTime;
import java.util.List;
import java.util.Objects;
import java.util.UUID;

public record RoomSubmitEventMessage(
UUID eventId, // 카프카 재전달 시 중복 처리 방지용 식별자 (Redis Dedup 키)
UUID userId,
UUID roomId,
List<ProblemAnswerRequest> answers,
OffsetDateTime submittedAt
) {

public RoomSubmitEventMessage {
Objects.requireNonNull(eventId, "eventId는 null일 수 없습니다");
Objects.requireNonNull(userId, "userId는 null일 수 없습니다");
Objects.requireNonNull(roomId, "roomId는 null일 수 없습니다");
Objects.requireNonNull(answers, "answers는 null일 수 없습니다");
Objects.requireNonNull(submittedAt, "submittedAt은 null일 수 없습니다");
answers = List.copyOf(answers);
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

public static RoomSubmitEventMessage of(UUID userId, UUID roomId, List<ProblemAnswerRequest> answers) {
return new RoomSubmitEventMessage(UUID.randomUUID(), userId, roomId, answers, OffsetDateTime.now());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,8 @@ public enum RoomErrorCode implements ErrorCode {
REPORT_GENERATION_FAILED(4008, "REPORT_GENERATION_FAILED", HttpStatus.INTERNAL_SERVER_ERROR, "리포트 PDF 문서 생성 중 서버 내부 오류가 발생했습니다."),
DUPLICATE_ANSWER_SUBMITTED(4009, "DUPLICATE_ANSWER_SUBMITTED", HttpStatus.BAD_REQUEST, "동일한 문제에 대한 중복 답안 제출은 허용되지 않습니다."),
REPORT_NOT_READY(4010, "REPORT_NOT_READY", HttpStatus.BAD_REQUEST, "아직 채점이 완료되지 않아 리포트가 준비되지 않았습니다."),
AI_GRADING_IN_PROGRESS(4011, "AI_GRADING_IN_PROGRESS", HttpStatus.BAD_REQUEST, "현재 해당 시험방의 AI 채점이 진행 중입니다.");
AI_GRADING_IN_PROGRESS(4011, "AI_GRADING_IN_PROGRESS", HttpStatus.BAD_REQUEST, "현재 해당 시험방의 AI 채점이 진행 중입니다."),
LOCK_ACQUISITION_FAILED(4012, "LOCK_ACQUISITION_FAILED", HttpStatus.SERVICE_UNAVAILABLE, "동시 요청 처리를 위한 락 획득에 실패했습니다. 잠시 후 다시 시도해주세요.");

private final int numeric;
private final String errorKey;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
package com.momogo.core.domain.room.kafka.producer;

import com.momogo.core.common.config.KafkaTopics;
import com.momogo.core.common.exception.BusinessException;
import com.momogo.core.common.exception.GlobalErrorCode;
import com.momogo.core.domain.room.event.RoomSubmitEventMessage;
import java.util.concurrent.TimeUnit;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Component;

/*
* 답안 제출 이벤트를 room-submit-events 토픽에 발행하는 전담 클래스.
* 카프카 브로커의 발행 확정(ACK)을 타임아웃 내에 동기적으로 대기하여 내구성을 보장하며,
* 실제 DB 저장은 momogo-api의 RoomSubmitKafkaConsumer가 비동기로 처리한다.
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class RoomSubmitKafkaProducer {

private final KafkaTemplate<String, Object> kafkaTemplate;

public void send(RoomSubmitEventMessage message) {
try {
SendResult<String, Object> result = kafkaTemplate.send(
KafkaTopics.ROOM_SUBMIT_EVENTS,
message.roomId().toString(),
message
).get(3, TimeUnit.SECONDS);

log.info("[RoomSubmitKafkaProducer] 카프카 메시지 전송 성공 - eventId: {}, offset: {}",
message.eventId(), result.getRecordMetadata().offset());
} catch (Exception ex) {
log.error("[RoomSubmitKafkaProducer] 카프카 메시지 전송 실패 - eventId: {}, roomId: {}, topic: {}",
message.eventId(), message.roomId(), KafkaTopics.ROOM_SUBMIT_EVENTS, ex);
if (ex instanceof InterruptedException) {
Thread.currentThread().interrupt();
}
throw new BusinessException(GlobalErrorCode.INTERNAL_SERVER_ERROR, "답안 제출 전송 실패");
}
}
Comment thread
SungHuii marked this conversation as resolved.
}
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.time.Duration;
import java.time.OffsetDateTime;
import java.util.ArrayList;
import java.util.Collections;
Expand All @@ -63,8 +64,17 @@
import java.util.stream.Collectors;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import com.momogo.core.common.exception.AuthErrorCode;
import com.momogo.core.common.exception.GlobalErrorCode;
import com.momogo.core.domain.room.event.RoomSubmitEventMessage;
import com.momogo.core.domain.room.kafka.producer.RoomSubmitKafkaProducer;
import java.util.concurrent.TimeUnit;
import org.redisson.api.RLock;
import org.redisson.api.RedissonClient;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Propagation;
import org.springframework.transaction.annotation.Transactional;

@Slf4j
Expand All @@ -84,6 +94,13 @@ public class RoomServiceImpl implements RoomService{
private final RoomMapper roomMapper;
private final ApplicationEventPublisher eventPublisher;
private final AiGradingProducer aiGradingProducer;
private final RedissonClient redissonClient;
private final RoomSubmitKafkaProducer roomSubmitKafkaProducer;
private final RedisTemplate<String, Object> redisTemplate;

private static final String SUBMIT_LOCK_PREFIX = "lock:room:submit:";
private static final String SUBMIT_CLAIM_PREFIX = "room:submit:claimed:";
private static final Duration SUBMIT_CLAIM_TTL = Duration.ofDays(7);

// 검증된 타겟 객체들을 묶기 위한 private record
private record ValidatedRoomTarget(Space space, List<User> targetUsers) {}
Expand Down Expand Up @@ -236,67 +253,84 @@ public List<RoomProblemResponse> getRoomProblems(UUID userId, UUID roomId) {
}

@Override
@Transactional
@Transactional(propagation = Propagation.NOT_SUPPORTED)
public void submitRoomAnswer(UUID userId, UUID roomId, RoomAnswerSubmitRequest request) {
log.info("[RoomService] 시험 답안 제출 시작 - userId: {}, roomId: {}", userId, roomId);
log.info("[RoomService] 시험 답안 제출 시도 - userId: {}, roomId: {}", userId, roomId);

// 요청 리스트 내 동일한 문제 ID가 중복 유입되는지 검증
// 1. 요청 리스트 내 동일한 문제 ID가 중복 유입되는지 검증
long uniqueProblemCount = request.answers().stream()
.map(ProblemAnswerRequest::roomProblemId)
.distinct().count();
if (uniqueProblemCount != request.answers().size()) {
throw new BusinessException(RoomErrorCode.DUPLICATE_ANSWER_SUBMITTED);
}

// 방 존재 검증
Room room = roomRepository.findByIdForUpdate(roomId)
.orElseThrow(() -> new BusinessException(RoomErrorCode.ROOM_NOT_FOUND));
// 2. Redisson 분산 락 획득 (Redisson Watchdog 활용을 위해 leaseTime 생략)
String lockKey = SUBMIT_LOCK_PREFIX + roomId + ":" + userId;
RLock lock = redissonClient.getLock(lockKey);

// 시간 검증 - 시험 시작 전 제출 시도 차단
OffsetDateTime now = OffsetDateTime.now();
if (now.isBefore(room.getTestStartAt())) {
throw new BusinessException(RoomErrorCode.INVALID_ACCESS_BEFORE_START);
}
try {
if (lock.tryLock(3, TimeUnit.SECONDS)) {
Comment thread
coderabbitai[bot] marked this conversation as resolved.
try {
// 3. 비관적 락(findByIdForUpdate) 대신 락이 없는 일반 조회 사용
Room room = findRoomOrThrow(roomId);

// 상태 검증 - 이미 최종적으로 끝난 시험이면 추가 제출 불가
if (Boolean.TRUE.equals(room.getIsEnded())) {
throw new BusinessException(RoomErrorCode.ALREADY_ENDED);
}
// 4. 시간 및 응시 자격 빠른 검증
OffsetDateTime now = OffsetDateTime.now();
if (now.isBefore(room.getTestStartAt())) {
throw new BusinessException(RoomErrorCode.INVALID_ACCESS_BEFORE_START);
}

// 응시 대상 유저 자격 검증
RoomUserId roomUserId = new RoomUserId(roomId, userId);
RoomUser roomUser = roomUserRepository.findByIdForUpdate(roomUserId)
.orElseThrow(() -> new BusinessException(RoomErrorCode.NOT_ROOM_PARTICIPANT));
if (Boolean.TRUE.equals(room.getIsEnded())) {
throw new BusinessException(RoomErrorCode.ALREADY_ENDED);
}

// 요청 유저 획득
User user = userRepository.findById(userId)
.orElseThrow(() -> new BusinessException(SpaceErrorCode.SPACE_USER_NOT_FOUND));
RoomUserId roomUserId = new RoomUserId(roomId, userId);
RoomUser roomUser = roomUserRepository.findById(roomUserId)
.orElseThrow(() -> new BusinessException(RoomErrorCode.NOT_ROOM_PARTICIPANT));

List<UUID> problemIds = request.answers().stream().map(ProblemAnswerRequest::roomProblemId).toList();
Map<UUID, RoomProblem> problemsById = roomProblemRepository.findAllById(problemIds).stream()
.collect(Collectors.toMap(RoomProblem::getId, Function.identity()));
if (Boolean.TRUE.equals(roomUser.getIsAttended())) {
throw new BusinessException(RoomErrorCode.ALREADY_ENDED);
}

if (Boolean.TRUE.equals(roomUser.getIsAttended())) {
throw new BusinessException(RoomErrorCode.ALREADY_ENDED);
}
// 4-1. 비동기 저장 완료 전 중복 제출을 즉시 차단하는 Redis Claim 마커 (Atomic SETNX)
String claimKey = SUBMIT_CLAIM_PREFIX + roomId + ":" + userId;
Boolean claimed = redisTemplate.opsForValue().setIfAbsent(claimKey, "1", SUBMIT_CLAIM_TTL);
if (!Boolean.TRUE.equals(claimed)) {
log.warn("[RoomService] 이미 답안 제출 요청이 처리 중이거나 완료된 유저 - userId: {}, roomId: {}", userId, roomId);
throw new BusinessException(RoomErrorCode.ALREADY_ENDED);
}

// 5. 문제 존재 및 해당 시험방 소속 검증
List<UUID> problemIds = request.answers().stream()
.map(ProblemAnswerRequest::roomProblemId)
.toList();

// 각 문제 답안 일괄 매핑, 저장
List<UserRoomAnswer> userRoomAnswers = request.answers().stream()
.map(ans -> {
RoomProblem problem = problemsById.get(ans.roomProblemId());
if (problem == null || !problem.getRoom().getId().equals(roomId)) {
long validProblemCount = roomProblemRepository.findAllById(problemIds).stream()
.filter(p -> p.getRoom().getId().equals(roomId))
.count();

if (validProblemCount != problemIds.size()) {
throw new BusinessException(RoomErrorCode.PROBLEM_NOT_FOUND);
}
return UserRoomAnswer.of(user, problem, ans.userAnswer(), null);
})
.toList();
userRoomAnswerRepository.saveAll(userRoomAnswers);

// 응시 완료 상태 업데이트
roomUser.attend();
// 6. 동기 DB 저장 대신 Kafka 전송 후 < 20ms 즉시 응답 반환
roomSubmitKafkaProducer.send(RoomSubmitEventMessage.of(userId, roomId, request.answers()));

log.info("[RoomService] 시험 답안 제출 완료 - userId: {}, roomId: {}, 제출 문항 수: {}",
userId, roomId, userRoomAnswers.size());
} finally {
if (lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
} else {
log.warn("[RoomService] 답안 제출 분산 락 획득 타임아웃 - userId: {}, roomId: {}", userId, roomId);
throw new BusinessException(RoomErrorCode.LOCK_ACQUISITION_FAILED);
}
Comment thread
SungHuii marked this conversation as resolved.
} catch (InterruptedException e) {
log.error("[RoomService] 답안 제출 대기 중 인터럽트 발생 - userId: {}, roomId: {}", userId, roomId, e);
Thread.currentThread().interrupt();
throw new BusinessException(GlobalErrorCode.INTERNAL_SERVER_ERROR, "답안 제출 대기 중 인터럽트가 발생했습니다.");
}
}

@Override
Expand Down