diff --git a/src/main/java/com/dongsoop/dongsoop/blinddate/entity/SessionInfo.java b/src/main/java/com/dongsoop/dongsoop/blinddate/entity/SessionInfo.java index e5ff1bcd1..e27e01c37 100644 --- a/src/main/java/com/dongsoop/dongsoop/blinddate/entity/SessionInfo.java +++ b/src/main/java/com/dongsoop/dongsoop/blinddate/entity/SessionInfo.java @@ -10,7 +10,7 @@ public class SessionInfo { private final String sessionId; private final LocalDateTime createdAt; - private SessionState state; + private volatile SessionState state; public static SessionInfo create() { return SessionInfo.builder() diff --git a/src/main/java/com/dongsoop/dongsoop/blinddate/executor/BlindDateEventQueue.java b/src/main/java/com/dongsoop/dongsoop/blinddate/executor/BlindDateEventQueue.java new file mode 100644 index 000000000..881ce34ce --- /dev/null +++ b/src/main/java/com/dongsoop/dongsoop/blinddate/executor/BlindDateEventQueue.java @@ -0,0 +1,65 @@ +package com.dongsoop.dongsoop.blinddate.executor; + +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; + +/** + * 과팅(BlindDate) 상태 변경 이벤트를 순서대로 처리하는 단일 스레드 큐. + *

+ * 입장/재접속/세션 시작 판단/퇴장/인원 브로드캐스트 등 참가자·세션 상태를 건드리는 작업은 + * 전부 이 큐를 통해 제출되어 하나의 스레드에서 순서대로 처리된다. 처리를 담당하는 스레드가 + * 항상 하나뿐이므로, 이 상태들에 대해서는 별도의 락이 필요 없다. + * (BlindDateMatchingLock, BlindDateMemberLock, BlindDateSessionLock 대체) + */ +@Slf4j +@Component +public class BlindDateEventQueue { + + private final ExecutorService executor = Executors.newSingleThreadExecutor(); + + /** + * 이벤트를 큐에 넣는다. 이미 큐에 있는 다른 이벤트들이 처리된 뒤 순서대로 실행된다. + *

+ * executor가 이미 종료된 상태에서 호출되어도(RejectedExecutionException 등) 호출자 스레드로 + * 예외가 전파되지 않도록 방어한다. + */ + public void submit(Runnable event) { + try { + executor.execute(() -> { + try { + event.run(); + } catch (Exception e) { + log.error("[BlindDate] Event processing failed", e); + } + }); + } catch (Exception e) { + log.error("[BlindDate] Failed to submit event to queue", e); + } + } + + /** + * 애플리케이션(빈) 종료 시 큐 정리. 과팅 라운드 종료(scheduleAutoClose)와는 다른 시점이다 — + * 한 번 호출되면 이 executor는 다시 못 쓰게 되므로, 다음 과팅 라운드를 위해 라운드 종료 + * 시점에는 절대 호출하면 안 된다. + */ + public void shutdown() { + executor.shutdownNow(); + } + + /** + * 지금까지 제출된 이벤트가 모두 처리될 때까지 대기한다. + *

+ * 큐가 단일 스레드로 순서대로(FIFO) 처리되므로, 트리비얼 작업을 제출해 그 작업이 끝날 때까지 + * 기다리면 그 이전에 제출된 모든 이벤트의 처리가 끝났음을 보장할 수 있다. 테스트에서 비동기 + * 처리 완료를 기다릴 때 사용한다. + */ + public void awaitIdle() { + try { + executor.submit(() -> null).get(); + } catch (Exception e) { + throw new IllegalStateException("[BlindDate] Failed to await queue idle", e); + } + } +} diff --git a/src/main/java/com/dongsoop/dongsoop/blinddate/handler/BlindDateConnectHandler.java b/src/main/java/com/dongsoop/dongsoop/blinddate/handler/BlindDateConnectHandler.java index efff2b39e..1540cd188 100644 --- a/src/main/java/com/dongsoop/dongsoop/blinddate/handler/BlindDateConnectHandler.java +++ b/src/main/java/com/dongsoop/dongsoop/blinddate/handler/BlindDateConnectHandler.java @@ -4,9 +4,7 @@ import com.dongsoop.dongsoop.blinddate.entity.ParticipantInfo; import com.dongsoop.dongsoop.blinddate.entity.SessionInfo; import com.dongsoop.dongsoop.blinddate.exception.SessionTerminatedException; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateMatchingLock; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateMemberLock; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateSessionLock; +import com.dongsoop.dongsoop.blinddate.executor.BlindDateEventQueue; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorage; import com.dongsoop.dongsoop.blinddate.repository.BlindDateSessionStorage; import com.dongsoop.dongsoop.blinddate.repository.BlindDateStorage; @@ -32,9 +30,8 @@ public class BlindDateConnectHandler { private final BlindDateSessionService sessionService; private final BlindDateSessionScheduler sessionScheduler; private final SimpMessagingTemplate messagingTemplate; - private final BlindDateMatchingLock blindDateMatchingLock; - private final BlindDateMemberLock blindDateMemberLock; - private final BlindDateSessionLock sessionLock; + private final BlindDateEventQueue eventQueue; + private final BlindDateDisconnectHandler disconnectHandler; /** * 세션 참여 및 세션 id 반환 @@ -46,8 +43,24 @@ public void execute(String socketId, Long memberId, Map sessionA // 과팅 운영 중이 아닌 경우 종료 this.validateBlindDateAvailability(); + // 참가자/세션 상태를 건드리는 처리는 전부 큐에서 순서대로 처리 + eventQueue.submit(() -> handle(socketId, memberId, sessionAttributes)); + } + + private void handle(String socketId, Long memberId, Map sessionAttributes) { // 이미 참여 중인 경우 소켓만 추가 후 종료 - String existingSessionId = this.tryHandleReconnection(socketId, memberId); + String existingSessionId; + try { + existingSessionId = this.tryHandleReconnection(socketId, memberId); + } catch (SessionTerminatedException e) { + // 재접속하려는 세션이 이미 종료된 경우: 참가자 기록은 재입장 방지를 위해 그대로 두고 + // 클라이언트에만 알림 (BlindDateSessionSchedulerImpl.finalizeSession 참고) + log.info("[BlindDate] Reconnect target session already terminated: memberId={}", memberId); + sessionAttributes.remove("sessionId"); + sendSessionTerminatedEvent(memberId); + return; + } + if (existingSessionId != null) { sessionAttributes.put("sessionId", existingSessionId); return; @@ -63,26 +76,18 @@ public void execute(String socketId, Long memberId, Map sessionA String sessionId = joinResult.sessionId(); - // 세션 시작 시 연결 해제 및 추가 연결을 막기 위한 세션 락 - this.sessionLock.lockBySessionId(sessionId); - - try { - // 마지막 참여자인지 검증 후 과팅 세션 시작 시도 - if (tryStart(sessionId)) { - // 마지막으로 입장한 사용자의 소켓 수신을 위해 현재 스레드를 종료하고 새 스레드에서 처리 - new Thread(() -> sessionScheduler.start(sessionId)).start(); - return; - } - } finally { - // 세션 락 해제 - this.sessionLock.unlockBySessionId(sessionId); + // 마지막 참여자인지 검증 후 과팅 세션 시작 시도 + if (tryStart(sessionId)) { + // 마지막으로 입장한 사용자의 소켓 수신을 위해 현재 스레드를 종료하고 새 스레드에서 처리 + new Thread(() -> sessionScheduler.start(sessionId)).start(); + return; } // 마지막 참여자가 아닌 경우 인원 업데이트 브로드캐스트 blindDateService.broadcastJoinedCount(joinResult.sessionId(), joinResult.currentCount()); } - private synchronized boolean tryStart(String sessionId) { + private boolean tryStart(String sessionId) { // 마지막 참여자인 경우 세션 시작 if (sessionService.isSessionFull(sessionId)) { if (!sessionStorage.isWaiting(sessionId)) { @@ -99,49 +104,50 @@ private synchronized boolean tryStart(String sessionId) { } private BlindDateJoinResult join(String socketId, Long memberId, Map sessionAttributes) { - // 처음 입장 시 포인터 할당을 위해 매칭 획득 - blindDateMatchingLock.lock(); - - BlindDateJoinResult joinResult; + String sessionId; + ParticipantInfo participant; try { - // 과팅 세션 할당 (Pointer 기반, Lock으로 동시성 보장) - String sessionId = assignSession(); + // 과팅 세션 할당 (Pointer 기반, 큐에서 순서대로 처리되므로 동시성 보장) + sessionId = assignSession(); // 과팅 세션 id 세션 속성에 저장 sessionAttributes.put("sessionId", sessionId); // 참여 정보 추가 (assignSession에서 편입 가능한 과팅 세션 여부를 확인했기에 바로 저장) - ParticipantInfo participant = participantStorage.addParticipant(sessionId, memberId, socketId); - - // 과팅 세션 편입 후 참가자 수 조회 - List participantInfos = participantStorage.findAllBySessionId(sessionId); - int currentCount = participantInfos.size(); - int maxCount = blindDateStorage.getMaxSessionMemberCount(); - - joinResult = new BlindDateJoinResult(participant, sessionId, currentCount, maxCount); + participant = participantStorage.addParticipant(sessionId, memberId, socketId); } catch (Exception e) { - // 입장 과정에서 오류 발생 시 회원 제거 + // addParticipant는 compute() 기반이라 예외 발생 시 참가자 맵에 아무 것도 반영되지 않는다. + // 따라서 여기서 회원을 제거하면, 다른 세션에 이미 정상 등록된 참가자 정보를 잘못 지울 수 있다. log.error("[BlindDate] Exception from enter process: memberId={}", memberId, e); - this.participantStorage.removeParticipant(memberId); sessionAttributes.remove("sessionId"); return null; - } finally { - // 회원 편입 후 회원 락 해제 - blindDateMatchingLock.unlock(); } try { + // 과팅 세션 편입 후 참가자 수 조회 + List participantInfos = participantStorage.findAllBySessionId(sessionId); + int currentCount = participantInfos.size(); + int maxCount = blindDateStorage.getMaxSessionMemberCount(); + + BlindDateJoinResult joinResult = new BlindDateJoinResult(participant, sessionId, currentCount, maxCount); + // 입장한 사용자에게 정보 전달 sendJoinEvent(joinResult); + + return joinResult; } catch (Exception e) { - // 입장한 사용자에게 정보 전달 실패 시 소켓 연결 해지로 보고 Disconnect에서 처리하도록 종료 - log.info("[BlindDate] Failed to send JOIN event, rolling back participant: memberId={}", memberId, e); + // 이 시점엔 addParticipant가 이미 성공해 참가자가 등록된 상태다. 등록 이후 단계에서 + // 실패하면(정원 조회 실패, 알림 전송 실패 등) 참가자가 고아로 남으므로, 방금 추가한 + // 소켓을 그대로 퇴장 처리(큐에 위임)해 되돌리고, 클라이언트에는 재시도를 요청한다. + log.error("[BlindDate] Post-registration failure, rolling back: memberId={}", memberId, e); + sessionAttributes.remove("sessionId"); + sendJoinFailedEvent(memberId); + disconnectHandler.execute(socketId, memberId, sessionId); + return null; } - - return joinResult; } /** @@ -161,32 +167,28 @@ private void validateBlindDateAvailability() { * @return 기존 세션 ID (재연결인 경우), null (첫 연결인 경우) */ private String tryHandleReconnection(String socketId, Long memberId) { - // 이미 참여중인 시나리오에 대해 안전한 소켓 추가를 위해 회원 락 획득 - blindDateMemberLock.lockByMemberId(memberId); - - try { - ParticipantInfo existingParticipant = participantStorage.getByMemberId(memberId); - // 첫 매칭인 경우 - if (existingParticipant == null) { - return null; - } - - String existingSessionId = existingParticipant.getSessionId(); + ParticipantInfo existingParticipant = participantStorage.getByMemberId(memberId); + // 첫 매칭인 경우 + if (existingParticipant == null) { + return null; + } - // 재연결된 세션이 존재하지 않는 경우 (세션 종료 후 재연결 시도 등) 예외 처리 - if (this.sessionStorage.getState(existingSessionId) == null) { - throw new SessionTerminatedException(); - } + String existingSessionId = existingParticipant.getSessionId(); - existingParticipant.addSocket(socketId); - return existingSessionId; - } finally { - blindDateMemberLock.unlockByMemberId(memberId); // 회원 락 해제 + // 재연결된 세션이 존재하지 않는 경우 (세션 종료 후 재연결 시도 등) 예외 처리 + if (this.sessionStorage.getState(existingSessionId) == null) { + throw new SessionTerminatedException(); } + + // addParticipant는 같은 세션이면 소켓 추가 + socketId->memberId 역인덱스 갱신까지 함께 처리한다. + // existingParticipant.addSocket()을 직접 호출하면 역인덱스가 안 갱신되어, 이 소켓이 + // 나중에 연결을 끊어도 removeSocket()이 찾지 못해 참가자가 정리되지 않는 문제가 있었다. + participantStorage.addParticipant(existingSessionId, memberId, socketId); + return existingSessionId; } /** - * 세션 할당 (Lock으로 동시성 보장) + * 세션 할당 (큐에서 순서대로 처리되어 동시성 보장) * * @return 할당된 세션 ID */ @@ -244,4 +246,42 @@ private void sendJoinEvent(BlindDateJoinResult joinResult) { throw e; } } + + /** + * 등록 이후 단계(정원 조회, 알림 전송 등) 실패로 입장을 되돌렸음을 클라이언트에 알리고 재시도를 요청 + * + * @param memberId 알림 대상 회원 id + */ + private void sendJoinFailedEvent(Long memberId) { + Map event = Map.of("state", "FAILED"); + + try { + messagingTemplate.convertAndSendToUser( + memberId.toString(), + "/queue/blinddate/join", + event + ); + } catch (Exception e) { + log.error("Failed to send JOIN_FAILED event: memberId={}", memberId, e); + } + } + + /** + * 재접속하려는 세션이 이미 종료되었음을 클라이언트에 알림 + * + * @param memberId 알림 대상 회원 id + */ + private void sendSessionTerminatedEvent(Long memberId) { + Map event = Map.of("state", "TERMINATED"); + + try { + messagingTemplate.convertAndSendToUser( + memberId.toString(), + "/queue/blinddate/join", + event + ); + } catch (Exception e) { + log.error("Failed to send SESSION_TERMINATED event: memberId={}", memberId, e); + } + } } diff --git a/src/main/java/com/dongsoop/dongsoop/blinddate/handler/BlindDateDisconnectHandler.java b/src/main/java/com/dongsoop/dongsoop/blinddate/handler/BlindDateDisconnectHandler.java index 4538b423f..1e255cb89 100644 --- a/src/main/java/com/dongsoop/dongsoop/blinddate/handler/BlindDateDisconnectHandler.java +++ b/src/main/java/com/dongsoop/dongsoop/blinddate/handler/BlindDateDisconnectHandler.java @@ -2,9 +2,7 @@ import com.dongsoop.dongsoop.blinddate.entity.ParticipantInfo; import com.dongsoop.dongsoop.blinddate.entity.SessionInfo.SessionState; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateMatchingLock; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateMemberLock; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateSessionLock; +import com.dongsoop.dongsoop.blinddate.executor.BlindDateEventQueue; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorage; import com.dongsoop.dongsoop.blinddate.repository.BlindDateSessionStorage; import com.dongsoop.dongsoop.blinddate.service.BlindDateService; @@ -21,38 +19,31 @@ public class BlindDateDisconnectHandler { private final BlindDateParticipantStorage participantStorage; private final BlindDateSessionStorage sessionStorage; private final BlindDateService blindDateService; - private final BlindDateMatchingLock matchingLock; - private final BlindDateMemberLock memberLock; - private final BlindDateSessionLock sessionLock; + private final BlindDateEventQueue eventQueue; public void execute(String socketId, Long memberId, String sessionId) { - // 연결 해제 중 세션 상태 변경을 막기 위한 세션 락 - this.sessionLock.lockBySessionId(sessionId); + // 참가자/세션 상태를 건드리는 처리는 전부 큐에서 순서대로 처리 + eventQueue.submit(() -> handle(socketId, memberId, sessionId)); + } - try { - // 종료된 세션에서 나가는 경우 별도 처리 안 함 - if (this.sessionStorage.getState(sessionId) == null) { - log.info("Session ID {} has been deleted", sessionId); - return; - } + private void handle(String socketId, Long memberId, String sessionId) { + // 종료된 세션에서 나가는 경우 별도 처리 안 함 + if (this.sessionStorage.getState(sessionId) == null) { + log.info("Session ID {} has been deleted", sessionId); + return; + } - // 참여자 정보에서 소켓 제거 시 남아있는 소켓이 없는지 - boolean isExit = this.removeSocketByParticipantInfo(socketId, memberId); + // 참여자 정보에서 소켓 제거 시 남아있는 소켓이 없는지 + boolean isExit = this.removeSocketByParticipantInfo(socketId, memberId); - // 사용자의 연결된 소켓이 없는 경우 퇴장 처리 - if (isExit && this.sessionStorage.isProcessing(sessionId)) { - // 참여중인 세션이 포인터 세션인 경우 회원도 제거 시도 - this.tryRemoveMember(memberId, sessionId); - } - } finally { - this.sessionLock.unlockBySessionId(sessionId); + // 사용자의 연결된 소켓이 없는 경우 퇴장 처리 + if (isExit && this.sessionStorage.isWaiting(sessionId)) { + // 참여중인 세션이 포인터 세션인 경우 회원도 제거 시도 + this.tryRemoveMember(memberId, sessionId); } } private boolean removeSocketByParticipantInfo(String socketId, Long memberId) { - // 소켓 제거 중 회원 소켓이 추가되지 않도록 처리 - memberLock.lockByMemberId(memberId); - try { // 소켓만 제거 (모든 소켓이 제거되면 자동으로 참여자도 제거됨) // 제거 후 모든 소켓을 제거했는지 여부 반환 @@ -60,35 +51,24 @@ private boolean removeSocketByParticipantInfo(String socketId, Long memberId) { } catch (IllegalArgumentException e) { log.warn("Socket already removed or not found: socketId={}, memberId={}", socketId, memberId); return false; - } finally { - // 락 해제 - memberLock.unlockByMemberId(memberId); } } private void tryRemoveMember(Long memberId, String sessionId) { - // 퇴장하려는 세션이 포인터 세션일 수 있기에 매칭 락 - matchingLock.lock(); + // 세션 상태 확인 + SessionState state = sessionStorage.getState(sessionId); - try { - // 세션 상태 확인 - SessionState state = sessionStorage.getState(sessionId); - - // PROCESSING 상태면 포인터가 아니며, 퇴장 처리 안 함 - if (state != SessionState.WAITING) { - return; - } + // PROCESSING 상태면 포인터가 아니며, 퇴장 처리 안 함 + if (state != SessionState.WAITING) { + return; + } - // WAITING 상태일 때만 퇴장 처리 - participantStorage.removeParticipant(memberId); + // WAITING 상태일 때만 퇴장 처리 + participantStorage.removeParticipant(memberId); - List participantInfos = participantStorage.findAllBySessionId(sessionId); + List participantInfos = participantStorage.findAllBySessionId(sessionId); - // 인원 업데이트 브로드캐스트 - blindDateService.broadcastJoinedCount(sessionId, participantInfos.size()); - } finally { - // 락 해제 - matchingLock.unlock(); - } + // 인원 업데이트 브로드캐스트 + blindDateService.broadcastJoinedCount(sessionId, participantInfos.size()); } } diff --git a/src/main/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateMatchingLock.java b/src/main/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateMatchingLock.java deleted file mode 100644 index a688c4958..000000000 --- a/src/main/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateMatchingLock.java +++ /dev/null @@ -1,23 +0,0 @@ -package com.dongsoop.dongsoop.blinddate.lock; - -import java.util.concurrent.locks.Lock; -import java.util.concurrent.locks.ReentrantLock; -import org.springframework.stereotype.Component; - -@Component -public class BlindDateMatchingLock { - - // 매칭 락 - private final Lock matchingLocks = new ReentrantLock(); - - /** - * 매칭 락 획득 - */ - public void lock() { - this.matchingLocks.lock(); - } - - public void unlock() { - this.matchingLocks.unlock(); - } -} diff --git a/src/main/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateMemberLock.java b/src/main/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateMemberLock.java deleted file mode 100644 index d2c04a2a5..000000000 --- a/src/main/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateMemberLock.java +++ /dev/null @@ -1,41 +0,0 @@ -package com.dongsoop.dongsoop.blinddate.lock; - -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.locks.Lock; -import java.util.concurrent.locks.ReentrantLock; -import org.springframework.stereotype.Component; - -@Component -public class BlindDateMemberLock { - - // 회원 락 - private final Map memberLocks = new ConcurrentHashMap<>(); - - /** - * 회원별 락 획득 - * - * @param memberId 락을 획득할 회원 id - */ - public void lockByMemberId(Long memberId) { - Lock lock = memberLocks.computeIfAbsent(memberId, k -> new ReentrantLock()); - lock.lock(); - } - - /** - * 회원 락 반납 - * - * @param memberId 락을 반납할 회원 id - */ - public void unlockByMemberId(Long memberId) { - Lock lock = memberLocks.computeIfAbsent(memberId, k -> new ReentrantLock()); - lock.unlock(); - } - - /** - * 세션 종료 후 락 클리어 - */ - public void clear() { - this.memberLocks.clear(); - } -} diff --git a/src/main/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateSessionLock.java b/src/main/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateSessionLock.java deleted file mode 100644 index b83a0f914..000000000 --- a/src/main/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateSessionLock.java +++ /dev/null @@ -1,41 +0,0 @@ -package com.dongsoop.dongsoop.blinddate.lock; - -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.locks.Lock; -import java.util.concurrent.locks.ReentrantLock; -import org.springframework.stereotype.Component; - -@Component -public class BlindDateSessionLock { - - // 세션 락 - private final Map sessionLocks = new ConcurrentHashMap<>(); - - /** - * 세션별 락 획득 - * - * @param sessionId 락을 획득할 세션 id - */ - public void lockBySessionId(String sessionId) { - Lock lock = this.sessionLocks.computeIfAbsent(sessionId, k -> new ReentrantLock()); - lock.lock(); - } - - /** - * 세션 락 반납 - * - * @param sessionId 락을 반납할 세션 id - */ - public void unlockBySessionId(String sessionId) { - Lock lock = sessionLocks.computeIfAbsent(sessionId, k -> new ReentrantLock()); - lock.unlock(); - } - - /** - * 세션 종료 후 락 클리어 - */ - public void clear() { - this.sessionLocks.clear(); - } -} diff --git a/src/main/java/com/dongsoop/dongsoop/blinddate/repository/BlindDateParticipantStorageImpl.java b/src/main/java/com/dongsoop/dongsoop/blinddate/repository/BlindDateParticipantStorageImpl.java index 59116af8a..789ad13e8 100644 --- a/src/main/java/com/dongsoop/dongsoop/blinddate/repository/BlindDateParticipantStorageImpl.java +++ b/src/main/java/com/dongsoop/dongsoop/blinddate/repository/BlindDateParticipantStorageImpl.java @@ -18,6 +18,9 @@ public class BlindDateParticipantStorageImpl implements BlindDateParticipantStor // memberId -> ParticipantInfo (한 사용자당 1개, 여러 소켓 보유 가능) private final Map participants = new ConcurrentHashMap<>(); + // socketId -> memberId (socketId로 참여자를 O(1)에 찾기 위한 역방향 인덱스) + private final Map socketIdToMemberId = new ConcurrentHashMap<>(); + // sessionId -> 익명 번호 카운터 private final Map nameCounters = new ConcurrentHashMap<>(); @@ -29,55 +32,59 @@ public class BlindDateParticipantStorageImpl implements BlindDateParticipantStor /** * 참여자 추가 또는 소켓 추가 + *

+ * memberId 기준으로 ConcurrentHashMap#compute를 사용해 조회와 반영을 원자적으로 묶는다. (외부에서 별도 락을 + * 잡고 호출하지 않아도 이 메서드 자체로 스레드 안전함) * * @param sessionId 참여하려는 세션 id * @param memberId 참여 주체 회원 id * @param socketId 참여 주체 소켓 id */ - public synchronized ParticipantInfo addParticipant(String sessionId, Long memberId, String socketId) { - // 이미 참여 중인지 확인 - ParticipantInfo existing = participants.get(memberId); - - // 처음 참여하는 경우 - if (existing == null) { - // 익명 이름 생성 (synchronized 메서드이므로 완전히 순차적으로 처리됨) - AtomicInteger atomicCounter = nameCounters.computeIfAbsent(sessionId, k -> new AtomicInteger(1)); - int counter = atomicCounter.getAndIncrement(); - String anonymousName = "익명" + counter; - - ParticipantInfo participant = ParticipantInfo.create(sessionId, memberId, socketId, anonymousName); - participants.put(memberId, participant); - - log.info("[BlindDate] Participant added: sessionId={}, memberId={}, socketId={}, name={}", - sessionId, memberId, socketId, anonymousName); - - return participant; - } - - // 참여중인 경우 - // 같은 세션이면 소켓만 추가 - if (existing.getSessionId().equals(sessionId)) { - existing.addSocket(socketId); - log.info("[BlindDate] Socket added to existing participant: memberId={}, socketId={}, totalSockets={}", - memberId, socketId, existing.getSocketIds().size()); - return existing; - } - - // 다른 세션에 이미 참여 중 - throw new IllegalStateException( - String.format("[BlindDate] Member %d already in session %s, cannot join session %s", - memberId, existing.getSessionId(), sessionId)); + public ParticipantInfo addParticipant(String sessionId, Long memberId, String socketId) { + ParticipantInfo participant = participants.compute(memberId, (id, existing) -> { + // 처음 참여하는 경우 + if (existing == null) { + AtomicInteger atomicCounter = nameCounters.computeIfAbsent(sessionId, k -> new AtomicInteger(1)); + int counter = atomicCounter.getAndIncrement(); + String anonymousName = "익명" + counter; + + ParticipantInfo created = ParticipantInfo.create(sessionId, memberId, socketId, anonymousName); + + log.info("[BlindDate] Participant added: sessionId={}, memberId={}, socketId={}, name={}", + sessionId, memberId, socketId, anonymousName); + + return created; + } + + // 참여중인 경우 + // 같은 세션이면 소켓만 추가 + if (existing.getSessionId().equals(sessionId)) { + existing.addSocket(socketId); + log.info("[BlindDate] Socket added to existing participant: memberId={}, socketId={}, totalSockets={}", + memberId, socketId, existing.getSocketIds().size()); + + return existing; + } + + // 다른 세션에 이미 참여 중 + throw new IllegalStateException( + String.format("[BlindDate] Member %d already in session %s, cannot join session %s", + memberId, existing.getSessionId(), sessionId)); + }); + + // socketId -> memberId 인덱스 갱신 (compute가 예외 없이 끝난 경우에만 도달) + socketIdToMemberId.put(socketId, memberId); + + return participant; } /** * 소켓 제거 (연결 해제) 모든 소켓이 제거되면 참여자도 제거 */ public boolean removeSocket(String socketId) throws IllegalArgumentException { - // socketId를 가진 참여자 찾기 - ParticipantInfo participant = participants.values().stream() - .filter(p -> p.getSocketIds().contains(socketId)) - .findFirst() - .orElse(null); + // socketId -> memberId 인덱스로 O(1) 조회 + Long memberId = socketIdToMemberId.remove(socketId); + ParticipantInfo participant = memberId != null ? participants.get(memberId) : null; if (participant == null) { log.warn("[BlindDate] Participant not found for socketId: {}", socketId); @@ -117,10 +124,8 @@ public ParticipantInfo getByMemberId(Long memberId) { * 소켓 ID로 참여 정보 조회 */ public ParticipantInfo getBySocketId(String socketId) { - return participants.values().stream() - .filter(p -> p.getSocketIds().contains(socketId)) - .findFirst() - .orElse(null); + Long memberId = socketIdToMemberId.get(socketId); + return memberId != null ? participants.get(memberId) : null; } /** @@ -148,16 +153,15 @@ public Map getParticipantsIdAndName(String sessionId) { /** * 선택 기록 */ - public synchronized boolean recordChoice(String sessionId, Long choicerId, Long targetId) { + public boolean recordChoice(String sessionId, Long choicerId, Long targetId) { Map sessionChoices = choices.computeIfAbsent(sessionId, k -> new ConcurrentHashMap<>()); - // 이미 선택했는지 확인 - if (sessionChoices.containsKey(choicerId)) { + // 이미 선택했는지 확인 + 반영을 원자적으로 처리 (확인과 반영 사이 공백 제거) + if (sessionChoices.putIfAbsent(choicerId, targetId) != null) { log.warn("Already chosen: sessionId={}, choicerId={}", sessionId, choicerId); return false; } - sessionChoices.put(choicerId, targetId); log.info("Choice recorded: sessionId={}, choicerId={} -> targetId={}", sessionId, choicerId, targetId); // 매칭 확인 @@ -189,6 +193,7 @@ public boolean isMatched(String sessionId, Long memberId) { */ public synchronized void clear() { participants.clear(); + socketIdToMemberId.clear(); nameCounters.clear(); choices.clear(); matches.clear(); @@ -202,6 +207,7 @@ public synchronized void clear() { public synchronized void removeParticipant(Long memberId) { ParticipantInfo participant = participants.remove(memberId); if (participant != null) { + participant.getSocketIds().forEach(socketIdToMemberId::remove); log.info("Participant removed: memberId={}, sessionId={}", memberId, participant.getSessionId()); } } diff --git a/src/main/java/com/dongsoop/dongsoop/blinddate/repository/BlindDateSessionStorageImpl.java b/src/main/java/com/dongsoop/dongsoop/blinddate/repository/BlindDateSessionStorageImpl.java index 4b72f2903..909f522ea 100644 --- a/src/main/java/com/dongsoop/dongsoop/blinddate/repository/BlindDateSessionStorageImpl.java +++ b/src/main/java/com/dongsoop/dongsoop/blinddate/repository/BlindDateSessionStorageImpl.java @@ -2,20 +2,16 @@ import com.dongsoop.dongsoop.blinddate.entity.SessionInfo; import com.dongsoop.dongsoop.blinddate.entity.SessionInfo.SessionState; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateSessionLock; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; -import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Repository; @Slf4j @Repository -@RequiredArgsConstructor public class BlindDateSessionStorageImpl implements BlindDateSessionStorage { private final Map sessions = new ConcurrentHashMap<>(); - private final BlindDateSessionLock sessionLock; /** * 세션 생성 @@ -63,15 +59,10 @@ public void start(String sessionId) { */ @Override public void terminate(String sessionId) { - sessionLock.lockBySessionId(sessionId); - try { - SessionInfo session = sessions.get(sessionId); - if (session != null) { - this.sessions.remove(sessionId); // 종료된 세션 정보 제거 - log.info("[BlindDate] Session terminated: sessionId={}", sessionId); - } - } finally { - sessionLock.unlockBySessionId(sessionId); + SessionInfo session = sessions.get(sessionId); + if (session != null) { + this.sessions.remove(sessionId); // 종료된 세션 정보 제거 + log.info("[BlindDate] Session terminated: sessionId={}", sessionId); } } @@ -79,12 +70,12 @@ public void terminate(String sessionId) { * 세션 전체 삭제 */ @Override - public synchronized void clear() { + public void clear() { this.sessions.clear(); } @Override - public synchronized boolean isWaiting(String sessionId) { + public boolean isWaiting(String sessionId) { SessionInfo sessionInfo = this.sessions.get(sessionId); if (sessionInfo == null) { return false; @@ -94,7 +85,7 @@ public synchronized boolean isWaiting(String sessionId) { } @Override - public synchronized boolean isProcessing(String sessionId) { + public boolean isProcessing(String sessionId) { SessionInfo sessionInfo = this.sessions.get(sessionId); if (sessionInfo == null) { return false; diff --git a/src/main/java/com/dongsoop/dongsoop/blinddate/scheduler/BlindDateSessionSchedulerImpl.java b/src/main/java/com/dongsoop/dongsoop/blinddate/scheduler/BlindDateSessionSchedulerImpl.java index 3524d3044..008cee5c0 100644 --- a/src/main/java/com/dongsoop/dongsoop/blinddate/scheduler/BlindDateSessionSchedulerImpl.java +++ b/src/main/java/com/dongsoop/dongsoop/blinddate/scheduler/BlindDateSessionSchedulerImpl.java @@ -2,6 +2,7 @@ import com.dongsoop.dongsoop.blinddate.config.BlindDateMessageProvider; import com.dongsoop.dongsoop.blinddate.config.BlindDateTopic; +import com.dongsoop.dongsoop.blinddate.executor.BlindDateEventQueue; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorage; import com.dongsoop.dongsoop.blinddate.repository.BlindDateSessionStorage; import java.util.List; @@ -30,6 +31,7 @@ public class BlindDateSessionSchedulerImpl implements BlindDateSessionScheduler private final BlindDateMessageProvider messageProvider; private final SimpMessagingTemplate messagingTemplate; private final BlindDateTaskScheduler taskScheduler; + private final BlindDateEventQueue eventQueue; @Value("${blinddate.event-message-amount}") private int eventMessageAmount; @@ -68,10 +70,10 @@ public void start(String sessionId) { } catch (InterruptedException e) { log.error("Session interrupted: {}", sessionId, e); Thread.currentThread().interrupt(); - this.sessionStorage.terminate(sessionId); + eventQueue.submit(() -> sessionStorage.terminate(sessionId)); } catch (Exception e) { log.error("Error in session: {}", sessionId, e); - this.sessionStorage.terminate(sessionId); + eventQueue.submit(() -> sessionStorage.terminate(sessionId)); } } @@ -115,7 +117,7 @@ private void sendEventMessage(int index, String sessionId, List eventMes MESSAGE_WAITING_TIME); } catch (Exception e) { log.error("Error sending event message {} for session {}", index, sessionId, e); - sessionStorage.terminate(sessionId); + eventQueue.submit(() -> sessionStorage.terminate(sessionId)); } } @@ -136,7 +138,7 @@ private void scheduledNextEventMessage(int index, String sessionId, List } catch (Exception e) { // 위 코드에서 예외 발생 시 세션 종료 log.error("[BlindDate] Error in scheduled thaw/next for session {}", sessionId, e); - sessionStorage.terminate(sessionId); + eventQueue.submit(() -> sessionStorage.terminate(sessionId)); } } @@ -177,7 +179,7 @@ private void finalizeSession(String sessionId) { sendFailedToUnmatched(sessionId); // 세션 종료 - sessionStorage.terminate(sessionId); + eventQueue.submit(() -> sessionStorage.terminate(sessionId)); // 회원 정보는 재 접속 방지를 위해 제거하지 않음 // participantStorage.clearSession(sessionId); diff --git a/src/main/java/com/dongsoop/dongsoop/blinddate/service/BlindDateServiceImpl.java b/src/main/java/com/dongsoop/dongsoop/blinddate/service/BlindDateServiceImpl.java index 2c921c2e3..10434af8d 100644 --- a/src/main/java/com/dongsoop/dongsoop/blinddate/service/BlindDateServiceImpl.java +++ b/src/main/java/com/dongsoop/dongsoop/blinddate/service/BlindDateServiceImpl.java @@ -2,6 +2,7 @@ import com.dongsoop.dongsoop.blinddate.config.BlindDateTopic; import com.dongsoop.dongsoop.blinddate.dto.StartBlindDateRequest; +import com.dongsoop.dongsoop.blinddate.executor.BlindDateEventQueue; import com.dongsoop.dongsoop.blinddate.notification.BlindDateNotification; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorage; import com.dongsoop.dongsoop.blinddate.repository.BlindDateSessionStorage; @@ -29,6 +30,7 @@ public class BlindDateServiceImpl implements BlindDateService { private final BlindDateSessionStorage sessionStorage; private final SimpMessagingTemplate messagingTemplate; private final BlindDateTaskScheduler taskScheduler; + private final BlindDateEventQueue eventQueue; /** * 과팅 운영 상태 확인 @@ -78,11 +80,14 @@ private void scheduleAutoClose() { // 과팅 상태 종료 blindDateStorage.close(); - // 모든 세션 정보 삭제 - sessionStorage.clear(); + // 세션/참가자 데이터 정리는 큐를 통해 순서대로 처리 + eventQueue.submit(() -> { + // 모든 세션 정보 삭제 + sessionStorage.clear(); - // 모든 참가자 정보 삭제 - participantStorage.clear(); + // 모든 참가자 정보 삭제 + participantStorage.clear(); + }); // TaskScheduler 정리 taskScheduler.cleanupAllSessions(); diff --git a/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateConcurrencyTest.java b/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateConcurrencyTest.java index f55be9e6e..dd33e14e8 100644 --- a/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateConcurrencyTest.java +++ b/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateConcurrencyTest.java @@ -4,10 +4,9 @@ import static org.mockito.Mockito.mock; import com.dongsoop.dongsoop.blinddate.entity.ParticipantInfo; +import com.dongsoop.dongsoop.blinddate.executor.BlindDateEventQueue; import com.dongsoop.dongsoop.blinddate.handler.BlindDateConnectHandler; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateMatchingLock; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateMemberLock; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateSessionLock; +import com.dongsoop.dongsoop.blinddate.handler.BlindDateDisconnectHandler; import com.dongsoop.dongsoop.blinddate.notification.BlindDateNotification; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorage; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorageImpl; @@ -62,22 +61,18 @@ class BlindDateConcurrencyTest { private BlindDateStorage blindDateStorage; private BlindDateSessionStorage sessionStorage; private BlindDateParticipantStorage participantStorage; - private BlindDateMatchingLock matchingLock; - private BlindDateMemberLock memberLock; + private BlindDateEventQueue eventQueue; private BlindDateSessionScheduler sessionScheduler; private BlindDateTaskScheduler taskScheduler; private SimpMessagingTemplate messagingTemplate; - private BlindDateSessionLock sessionLock; @BeforeEach void setUp() { blindDateStorage = new BlindDateStorageImpl(); participantStorage = new BlindDateParticipantStorageImpl(); - sessionStorage = new BlindDateSessionStorageImpl(sessionLock); + sessionStorage = new BlindDateSessionStorageImpl(); - matchingLock = new BlindDateMatchingLock(); - memberLock = new BlindDateMemberLock(); - sessionLock = new BlindDateSessionLock(); + eventQueue = new BlindDateEventQueue(); messagingTemplate = mock(SimpMessagingTemplate.class); BlindDateNotification notification = mock(BlindDateNotification.class); @@ -90,7 +85,8 @@ void setUp() { notification, sessionStorage, messagingTemplate, - taskScheduler + taskScheduler, + eventQueue ); sessionService = new BlindDateSessionServiceImpl( @@ -98,6 +94,13 @@ void setUp() { blindDateStorage ); + BlindDateDisconnectHandler disconnectHandler = new BlindDateDisconnectHandler( + participantStorage, + sessionStorage, + blindDateService, + eventQueue + ); + connectHandler = new BlindDateConnectHandler( participantStorage, blindDateStorage, @@ -106,9 +109,8 @@ void setUp() { sessionService, sessionScheduler, messagingTemplate, - matchingLock, - memberLock, - sessionLock + eventQueue, + disconnectHandler ); } @@ -145,6 +147,7 @@ void concurrentFirstEntry_ShouldCreateOnlyOneSession() throws InterruptedExcepti try { Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-" + memberId, (long) memberId, sessionAttributes); + eventQueue.awaitIdle(); String sessionId = (String) sessionAttributes.get("sessionId"); if (sessionId != null) { sessionIds.add(sessionId); @@ -159,6 +162,7 @@ void concurrentFirstEntry_ShouldCreateOnlyOneSession() throws InterruptedExcepti } assertThat(latch.await(15, TimeUnit.SECONDS)).isTrue(); + eventQueue.awaitIdle(); executor.shutdown(); executor.awaitTermination(5, TimeUnit.SECONDS); @@ -186,6 +190,7 @@ void concurrentFullSession_ShouldCreateOnlyOneNewSession() throws InterruptedExc connectHandler.execute("socket-2", 2L, attr); connectHandler.execute("socket-3", 3L, attr); connectHandler.execute("socket-4", 4L, attr); + eventQueue.awaitIdle(); String firstSessionId = (String) attr.get("sessionId"); // 첫 번째 세션이 정확히 4명인지 확인 @@ -204,6 +209,7 @@ void concurrentFullSession_ShouldCreateOnlyOneNewSession() throws InterruptedExc try { Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-" + memberId, (long) memberId, sessionAttributes); + eventQueue.awaitIdle(); String sessionId = (String) sessionAttributes.get("sessionId"); if (sessionId != null) { newSessionIds.add(sessionId); @@ -218,6 +224,7 @@ void concurrentFullSession_ShouldCreateOnlyOneNewSession() throws InterruptedExc } assertThat(latch.await(15, TimeUnit.SECONDS)).isTrue(); + eventQueue.awaitIdle(); executor.shutdown(); executor.awaitTermination(5, TimeUnit.SECONDS); @@ -255,6 +262,7 @@ void concurrentPointerChange_ShouldBeConsistent() throws InterruptedException { Map attr = new HashMap<>(); connectHandler.execute("socket-1", 1L, attr); connectHandler.execute("socket-2", 2L, attr); + eventQueue.awaitIdle(); // 첫 번째 세션에 2명이 있는지 확인 String firstSession = (String) attr.get("sessionId"); @@ -273,6 +281,7 @@ void concurrentPointerChange_ShouldBeConsistent() throws InterruptedException { try { Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-" + memberId, (long) memberId, sessionAttributes); + eventQueue.awaitIdle(); String sessionId = (String) sessionAttributes.get("sessionId"); if (sessionId != null) { memberToSession.put((long) memberId, sessionId); @@ -287,6 +296,7 @@ void concurrentPointerChange_ShouldBeConsistent() throws InterruptedException { } assertThat(latch.await(15, TimeUnit.SECONDS)).isTrue(); + eventQueue.awaitIdle(); executor.shutdown(); executor.awaitTermination(5, TimeUnit.SECONDS); @@ -339,6 +349,7 @@ void concurrentExactCapacity_AllShouldSucceed() throws InterruptedException { try { Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-" + memberId, (long) memberId, sessionAttributes); + eventQueue.awaitIdle(); String sessionId = (String) sessionAttributes.get("sessionId"); if (sessionId != null) { sessionIds.add(sessionId); @@ -353,6 +364,7 @@ void concurrentExactCapacity_AllShouldSucceed() throws InterruptedException { } assertThat(latch.await(15, TimeUnit.SECONDS)).isTrue(); + eventQueue.awaitIdle(); executor.shutdown(); executor.awaitTermination(5, TimeUnit.SECONDS); @@ -389,6 +401,7 @@ void concurrentOverCapacity_ShouldCreateNewSession() throws InterruptedException try { Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-" + memberId, (long) memberId, sessionAttributes); + eventQueue.awaitIdle(); String sessionId = (String) sessionAttributes.get("sessionId"); if (sessionId != null) { memberToSession.put((long) memberId, sessionId); @@ -403,6 +416,7 @@ void concurrentOverCapacity_ShouldCreateNewSession() throws InterruptedException } assertThat(latch.await(15, TimeUnit.SECONDS)).isTrue(); + eventQueue.awaitIdle(); executor.shutdown(); executor.awaitTermination(5, TimeUnit.SECONDS); @@ -445,6 +459,7 @@ void concurrentLastSlot_OnlyOneShouldEnter() throws InterruptedException { connectHandler.execute("socket-2", 2L, attr); connectHandler.execute("socket-3", 3L, attr); connectHandler.execute("socket-4", 4L, attr); + eventQueue.awaitIdle(); String firstSessionId = (String) attr.get("sessionId"); // 4명이 정확히 입장했는지 확인 @@ -464,6 +479,7 @@ void concurrentLastSlot_OnlyOneShouldEnter() throws InterruptedException { try { Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-" + memberId, (long) memberId, sessionAttributes); + eventQueue.awaitIdle(); String sessionId = (String) sessionAttributes.get("sessionId"); if (sessionId != null) { memberToSession.put((long) memberId, sessionId); @@ -478,6 +494,7 @@ void concurrentLastSlot_OnlyOneShouldEnter() throws InterruptedException { } assertThat(latch.await(15, TimeUnit.SECONDS)).isTrue(); + eventQueue.awaitIdle(); executor.shutdown(); executor.awaitTermination(5, TimeUnit.SECONDS); @@ -531,6 +548,7 @@ void concurrentMultipleTabs_ShouldShareSession() throws InterruptedException { try { Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-" + socketNum, 1L, sessionAttributes); + eventQueue.awaitIdle(); String sessionId = (String) sessionAttributes.get("sessionId"); if (sessionId != null) { sessionIds.add(sessionId); @@ -545,6 +563,7 @@ void concurrentMultipleTabs_ShouldShareSession() throws InterruptedException { } assertThat(latch.await(15, TimeUnit.SECONDS)).isTrue(); + eventQueue.awaitIdle(); executor.shutdown(); executor.awaitTermination(5, TimeUnit.SECONDS); @@ -586,6 +605,7 @@ void concurrentAnonymousName_ShouldBeUnique() throws InterruptedException { try { Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-" + memberId, (long) memberId, sessionAttributes); + eventQueue.awaitIdle(); ParticipantInfo participant = participantStorage.getByMemberId((long) memberId); if (participant != null) { memberToName.put((long) memberId, participant.getAnonymousName()); @@ -600,6 +620,7 @@ void concurrentAnonymousName_ShouldBeUnique() throws InterruptedException { } assertThat(latch.await(20, TimeUnit.SECONDS)).isTrue(); + eventQueue.awaitIdle(); executor.shutdown(); executor.awaitTermination(5, TimeUnit.SECONDS); @@ -654,6 +675,7 @@ void stressTest_100Users() throws InterruptedException { try { Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-" + memberId, (long) memberId, sessionAttributes); + eventQueue.awaitIdle(); String sessionId = (String) sessionAttributes.get("sessionId"); if (sessionId != null) { memberToSession.put((long) memberId, sessionId); @@ -668,6 +690,7 @@ void stressTest_100Users() throws InterruptedException { } boolean completed = latch.await(30, TimeUnit.SECONDS); + eventQueue.awaitIdle(); executor.shutdown(); long duration = System.currentTimeMillis() - startTime; @@ -724,6 +747,7 @@ void edgeCase_NullPointer() throws InterruptedException { try { Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-" + memberId, (long) memberId, sessionAttributes); + eventQueue.awaitIdle(); memberToSession.put((long) memberId, (String) sessionAttributes.get("sessionId")); } finally { latch.countDown(); @@ -732,6 +756,7 @@ void edgeCase_NullPointer() throws InterruptedException { } assertThat(latch.await(20, TimeUnit.SECONDS)).isTrue(); + eventQueue.awaitIdle(); executor.shutdown(); Set sessions = new HashSet<>(memberToSession.values()); @@ -757,6 +782,7 @@ void edgeCase_CapacityOne() throws InterruptedException { try { Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-" + memberId, (long) memberId, sessionAttributes); + eventQueue.awaitIdle(); memberToSession.put((long) memberId, (String) sessionAttributes.get("sessionId")); } finally { latch.countDown(); @@ -765,6 +791,7 @@ void edgeCase_CapacityOne() throws InterruptedException { } assertThat(latch.await(20, TimeUnit.SECONDS)).isTrue(); + eventQueue.awaitIdle(); executor.shutdown(); assertThat(memberToSession).hasSize(userCount); diff --git a/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateIntegrationTest.java b/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateIntegrationTest.java index 9ba2eb5f6..94a7b9ea9 100644 --- a/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateIntegrationTest.java +++ b/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateIntegrationTest.java @@ -13,9 +13,7 @@ import com.dongsoop.dongsoop.blinddate.handler.BlindDateChoiceHandler; import com.dongsoop.dongsoop.blinddate.handler.BlindDateConnectHandler; import com.dongsoop.dongsoop.blinddate.handler.BlindDateDisconnectHandler; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateMatchingLock; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateMemberLock; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateSessionLock; +import com.dongsoop.dongsoop.blinddate.executor.BlindDateEventQueue; import com.dongsoop.dongsoop.blinddate.notification.BlindDateNotification; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorage; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorageImpl; @@ -69,22 +67,18 @@ class BlindDateIntegrationTest { private BlindDateDisconnectHandler disconnectHandler; private BlindDateSessionScheduler sessionScheduler; private BlindDateTaskScheduler taskScheduler; - private BlindDateMatchingLock matchingLock; - private BlindDateMemberLock memberLock; + private BlindDateEventQueue eventQueue; private SimpMessagingTemplate messagingTemplate; - private BlindDateSessionLock sessionLock; @BeforeEach void setUp() { // Repository 초기화 blindDateStorage = new BlindDateStorageImpl(); participantStorage = new BlindDateParticipantStorageImpl(); - sessionStorage = new BlindDateSessionStorageImpl(sessionLock); + sessionStorage = new BlindDateSessionStorageImpl(); - // Lock 초기화 - matchingLock = new BlindDateMatchingLock(); - memberLock = new BlindDateMemberLock(); - sessionLock = new BlindDateSessionLock(); + // 큐 초기화 + eventQueue = new BlindDateEventQueue(); // Mock 초기화 messagingTemplate = mock(SimpMessagingTemplate.class); @@ -114,7 +108,8 @@ void setUp() { notification, sessionStorage, messagingTemplate, - taskScheduler + taskScheduler, + eventQueue ); sessionService = new BlindDateSessionServiceImpl( @@ -123,6 +118,13 @@ void setUp() { ); // Handler 초기화 + disconnectHandler = new BlindDateDisconnectHandler( + participantStorage, + sessionStorage, + blindDateService, + eventQueue + ); + connectHandler = new BlindDateConnectHandler( participantStorage, blindDateStorage, @@ -131,9 +133,8 @@ void setUp() { sessionService, sessionScheduler, messagingTemplate, - matchingLock, - memberLock, - sessionLock + eventQueue, + disconnectHandler ); choiceHandler = new BlindDateChoiceHandler( @@ -141,15 +142,6 @@ void setUp() { messagingTemplate, chatRoomService ); - - disconnectHandler = new BlindDateDisconnectHandler( - participantStorage, - sessionStorage, - blindDateService, - matchingLock, - memberLock, - sessionLock - ); } @AfterEach @@ -180,6 +172,7 @@ void firstUser_CreatesNewSession() { // when Map sessionAttributes = new HashMap<>(); connectHandler.execute("socket-1", 1L, sessionAttributes); + eventQueue.awaitIdle(); String sessionId = (String) sessionAttributes.get("sessionId"); // then @@ -206,6 +199,7 @@ void multipleUsers_AssignedToSameSession() { connectHandler.execute("socket-1", 1L, attr1); connectHandler.execute("socket-2", 2L, attr2); connectHandler.execute("socket-3", 3L, attr3); + eventQueue.awaitIdle(); String session1 = (String) attr1.get("sessionId"); String session2 = (String) attr2.get("sessionId"); @@ -227,11 +221,13 @@ void fullSession_CreatesNewSession() { connectHandler.execute("socket-1", 1L, attr1); connectHandler.execute("socket-2", 2L, attr1); connectHandler.execute("socket-3", 3L, attr1); + eventQueue.awaitIdle(); String session1 = (String) attr1.get("sessionId"); // 4번째 사용자 Map attr2 = new HashMap<>(); connectHandler.execute("socket-4", 4L, attr2); + eventQueue.awaitIdle(); String session2 = (String) attr2.get("sessionId"); // then @@ -246,11 +242,13 @@ void reconnect_ReturnsToExistingSession() { blindDateStorage.start(5, LocalDateTime.now().plusHours(1)); Map attr1 = new HashMap<>(); connectHandler.execute("socket-1", 1L, attr1); + eventQueue.awaitIdle(); String session1 = (String) attr1.get("sessionId"); // when - 같은 memberId로 재연결 Map attr2 = new HashMap<>(); connectHandler.execute("socket-2", 1L, attr2); + eventQueue.awaitIdle(); String session2 = (String) attr2.get("sessionId"); // then @@ -285,6 +283,7 @@ void addParticipant_SavesAllInfo() { // when Map attr = new HashMap<>(); connectHandler.execute("socket-1", 1L, attr); + eventQueue.awaitIdle(); String sessionId = (String) attr.get("sessionId"); // then @@ -309,6 +308,7 @@ void anonymousNames_AssignedSequentially() { connectHandler.execute("socket-1", 1L, attr); connectHandler.execute("socket-2", 2L, attr); connectHandler.execute("socket-3", 3L, attr); + eventQueue.awaitIdle(); // then assertThat(participantStorage.getAnonymousName(1L)).isEqualTo("익명1"); @@ -332,6 +332,7 @@ void lastParticipant_PrepareToStartSession() throws InterruptedException { connectHandler.execute("socket-1", 1L, attr); connectHandler.execute("socket-2", 2L, attr); connectHandler.execute("socket-3", 3L, attr); + eventQueue.awaitIdle(); String sessionId = (String) attr.get("sessionId"); // 세션 시작은 비동기이므로 잠시 대기 @@ -353,6 +354,7 @@ void notFull_RemainsWaiting() throws InterruptedException { connectHandler.execute("socket-1", 1L, attr); connectHandler.execute("socket-2", 2L, attr); connectHandler.execute("socket-3", 3L, attr); + eventQueue.awaitIdle(); String sessionId = (String) attr.get("sessionId"); Thread.sleep(500); @@ -375,12 +377,14 @@ void disconnect_RemovesSocket() { Map attr = new HashMap<>(); connectHandler.execute("socket-1", 1L, attr); connectHandler.execute("socket-2", 2L, attr); + eventQueue.awaitIdle(); String sessionId = (String) attr.get("sessionId"); assertThat(getParticipantCount(sessionId)).isEqualTo(2); // when disconnectHandler.execute("socket-1", 1L, sessionId); + eventQueue.awaitIdle(); // then assertThat(participantStorage.getByMemberId(1L)).isNull(); @@ -409,6 +413,7 @@ void concurrent100Users_AccurateDistribution() throws Exception { try { Map attr = new HashMap<>(); connectHandler.execute("socket-" + memberId, memberId, attr); + eventQueue.awaitIdle(); String sessionId = (String) attr.get("sessionId"); if (sessionId != null) { sessions.add(sessionId); @@ -425,6 +430,7 @@ void concurrent100Users_AccurateDistribution() throws Exception { assertThat(latch.await(30, TimeUnit.SECONDS)).as("모든 스레드가 완료되어야 함").isTrue(); executor.shutdown(); executor.awaitTermination(10, TimeUnit.SECONDS); + eventQueue.awaitIdle(); if (!exceptions.isEmpty()) { log.error("=== Exceptions in 100 users test ==="); @@ -473,6 +479,7 @@ void concurrentEnter_PointerSynchronized() throws Exception { assertThat(latch.await(15, TimeUnit.SECONDS)).as("모든 스레드가 완료되어야 함").isTrue(); executor.shutdown(); executor.awaitTermination(10, TimeUnit.SECONDS); + eventQueue.awaitIdle(); if (!exceptions.isEmpty()) { log.error("=== Exceptions in pointer sync test ==="); @@ -523,6 +530,7 @@ void mutualChoice_Matches() { Map attr = new HashMap<>(); connectHandler.execute("socket-1", 1L, attr); connectHandler.execute("socket-2", 2L, attr); + eventQueue.awaitIdle(); String sessionId = (String) attr.get("sessionId"); // when @@ -546,6 +554,7 @@ void oneWayChoice_NoMatch() { Map attr = new HashMap<>(); connectHandler.execute("socket-1", 1L, attr); connectHandler.execute("socket-2", 2L, attr); + eventQueue.awaitIdle(); String sessionId = (String) attr.get("sessionId"); // when @@ -568,6 +577,7 @@ void triangleChoice_AllFail() { connectHandler.execute("socket-1", 1L, attr); connectHandler.execute("socket-2", 2L, attr); connectHandler.execute("socket-3", 3L, attr); + eventQueue.awaitIdle(); String sessionId = (String) attr.get("sessionId"); // when - 1→2, 2→3, 3→1 @@ -590,6 +600,7 @@ void partialMatching() { for (long i = 2; i <= 5; i++) { connectHandler.execute("socket-" + i, i, attr); } + eventQueue.awaitIdle(); String sessionId = (String) attr.get("sessionId"); // when - 1↔2, 3↔4, 5 혼자 @@ -624,6 +635,7 @@ void completeFlow_StartToEnd() throws InterruptedException { connectHandler.execute("socket-1", 1L, attr); connectHandler.execute("socket-2", 2L, attr); connectHandler.execute("socket-3", 3L, attr); + eventQueue.awaitIdle(); String sessionId = (String) attr.get("sessionId"); Thread.sleep(500); diff --git a/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateLockVerificationTest.java b/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateLockVerificationTest.java deleted file mode 100644 index 2933c83f8..000000000 --- a/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateLockVerificationTest.java +++ /dev/null @@ -1,254 +0,0 @@ -package com.dongsoop.dongsoop.blinddate; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.assertj.core.api.Assertions.assertThatThrownBy; - -import com.dongsoop.dongsoop.blinddate.lock.BlindDateMemberLock; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateSessionLock; -import java.util.ArrayList; -import java.util.Collections; -import java.util.List; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicInteger; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.DisplayName; -import org.junit.jupiter.api.RepeatedTest; -import org.junit.jupiter.api.Test; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * Lock 메커니즘 검증 테스트 - *

- * 이 테스트는 실제 동시성 문제를 발견하기 위한 테스트입니다. - */ -@DisplayName("Lock 메커니즘 검증 테스트") -class BlindDateLockVerificationTest { - - private static final Logger log = LoggerFactory.getLogger(BlindDateLockVerificationTest.class); - - private BlindDateSessionLock sessionLock; - private BlindDateMemberLock memberLock; - - @BeforeEach - void setUp() { - sessionLock = new BlindDateSessionLock(); - memberLock = new BlindDateMemberLock(); - } - - @RepeatedTest(20) - @DisplayName("세션 락 동시 획득 - 순차적으로 실행되어야 함") - void sessionLock_ConcurrentAccess_ShouldBeSequential() throws InterruptedException { - String sessionId = "test-session"; - int threadCount = 10; - ExecutorService executor = Executors.newFixedThreadPool(threadCount); - CountDownLatch startLatch = new CountDownLatch(1); - CountDownLatch endLatch = new CountDownLatch(threadCount); - - AtomicInteger counter = new AtomicInteger(0); - AtomicInteger maxConcurrent = new AtomicInteger(0); - AtomicInteger currentConcurrent = new AtomicInteger(0); - List exceptions = Collections.synchronizedList(new ArrayList<>()); - - for (int i = 0; i < threadCount; i++) { - final int threadNum = i; - executor.submit(() -> { - try { - startLatch.await(); // 모든 스레드가 동시에 시작하도록 - - sessionLock.lockBySessionId(sessionId); - try { - int concurrent = currentConcurrent.incrementAndGet(); - maxConcurrent.updateAndGet(max -> Math.max(max, concurrent)); - - // Critical section - int value = counter.get(); - Thread.sleep(10); // 경합 조건 유도 - counter.set(value + 1); - - currentConcurrent.decrementAndGet(); - } finally { - sessionLock.unlockBySessionId(sessionId); - } - } catch (Exception e) { - exceptions.add(e); - log.error("Exception in thread {}", threadNum, e); - } finally { - endLatch.countDown(); - } - }); - } - - startLatch.countDown(); // 모든 스레드 시작 - assertThat(endLatch.await(10, TimeUnit.SECONDS)).isTrue(); - executor.shutdown(); - executor.awaitTermination(5, TimeUnit.SECONDS); - - if (!exceptions.isEmpty()) { - log.error("=== Lock Test Exceptions ==="); - exceptions.forEach(e -> log.error("Exception: {}", e.getMessage(), e)); - } - - assertThat(exceptions).as("예외 없어야 함").isEmpty(); - assertThat(maxConcurrent.get()).as("최대 동시 실행 수는 1이어야 함").isEqualTo(1); - assertThat(counter.get()).as("카운터는 정확히 " + threadCount + "이어야 함").isEqualTo(threadCount); - } - - @Test - @DisplayName("Null sessionId로 unlock 시도 - IllegalMonitorStateException 또는 NullPointerException 발생 가능") - void sessionLock_UnlockNullSessionId_ShouldHandleGracefully() { - // 이 테스트는 현재 구현의 버그를 확인하는 테스트입니다. - // null을 전달하면 computeIfAbsent가 NullPointerException을 던지거나 - // 잘못된 lock을 unlock하려고 시도할 수 있습니다. - - // 이상적으로는 null 체크를 해야 합니다. - assertThatThrownBy(() -> { - sessionLock.unlockBySessionId(null); - }).isInstanceOfAny(NullPointerException.class, IllegalMonitorStateException.class); - } - - @Test - @DisplayName("Lock을 획득하지 않고 unlock 시도 - IllegalMonitorStateException") - void sessionLock_UnlockWithoutLock_ShouldThrowException() { - String sessionId = "test-session"; - - assertThatThrownBy(() -> { - sessionLock.unlockBySessionId(sessionId); - }).isInstanceOf(IllegalMonitorStateException.class); - } - - @Test - @DisplayName("중복 unlock 시도 - IllegalMonitorStateException") - void sessionLock_DoubleUnlock_ShouldThrowException() { - String sessionId = "test-session"; - - sessionLock.lockBySessionId(sessionId); - sessionLock.unlockBySessionId(sessionId); - - // 두 번째 unlock은 예외를 던져야 함 - assertThatThrownBy(() -> { - sessionLock.unlockBySessionId(sessionId); - }).isInstanceOf(IllegalMonitorStateException.class); - } - - @RepeatedTest(20) - @DisplayName("여러 세션의 락 동시 획득 - 각 세션은 독립적이어야 함") - void sessionLock_MultipleSessions_ShouldBeIndependent() throws InterruptedException { - int sessionCount = 5; - int threadsPerSession = 4; - int totalThreads = sessionCount * threadsPerSession; - - ExecutorService executor = Executors.newFixedThreadPool(totalThreads); - CountDownLatch startLatch = new CountDownLatch(1); - CountDownLatch endLatch = new CountDownLatch(totalThreads); - - AtomicInteger[] counters = new AtomicInteger[sessionCount]; - for (int i = 0; i < sessionCount; i++) { - counters[i] = new AtomicInteger(0); - } - - List exceptions = Collections.synchronizedList(new ArrayList<>()); - - for (int s = 0; s < sessionCount; s++) { - final int sessionNum = s; - final String sessionId = "session-" + s; - - for (int t = 0; t < threadsPerSession; t++) { - executor.submit(() -> { - try { - startLatch.await(); - - sessionLock.lockBySessionId(sessionId); - try { - int value = counters[sessionNum].get(); - Thread.sleep(5); - counters[sessionNum].set(value + 1); - } finally { - sessionLock.unlockBySessionId(sessionId); - } - } catch (Exception e) { - exceptions.add(e); - log.error("Exception in session {}", sessionNum, e); - } finally { - endLatch.countDown(); - } - }); - } - } - - startLatch.countDown(); - assertThat(endLatch.await(15, TimeUnit.SECONDS)).isTrue(); - executor.shutdown(); - executor.awaitTermination(5, TimeUnit.SECONDS); - - if (!exceptions.isEmpty()) { - log.error("=== Multiple Sessions Test Exceptions ==="); - exceptions.forEach(e -> log.error("Exception: {}", e.getMessage(), e)); - } - - assertThat(exceptions).isEmpty(); - for (int i = 0; i < sessionCount; i++) { - assertThat(counters[i].get()).as("Session " + i + " 카운터").isEqualTo(threadsPerSession); - } - } - - @RepeatedTest(20) - @DisplayName("회원 락 동시 획득 - 순차적으로 실행되어야 함") - void memberLock_ConcurrentAccess_ShouldBeSequential() throws InterruptedException { - Long memberId = 1L; - int threadCount = 10; - ExecutorService executor = Executors.newFixedThreadPool(threadCount); - CountDownLatch startLatch = new CountDownLatch(1); - CountDownLatch endLatch = new CountDownLatch(threadCount); - - AtomicInteger counter = new AtomicInteger(0); - AtomicInteger maxConcurrent = new AtomicInteger(0); - AtomicInteger currentConcurrent = new AtomicInteger(0); - List exceptions = Collections.synchronizedList(new ArrayList<>()); - - for (int i = 0; i < threadCount; i++) { - final int threadNum = i; - executor.submit(() -> { - try { - startLatch.await(); - - memberLock.lockByMemberId(memberId); - try { - int concurrent = currentConcurrent.incrementAndGet(); - maxConcurrent.updateAndGet(max -> Math.max(max, concurrent)); - - int value = counter.get(); - Thread.sleep(10); - counter.set(value + 1); - - currentConcurrent.decrementAndGet(); - } finally { - memberLock.unlockByMemberId(memberId); - } - } catch (Exception e) { - exceptions.add(e); - log.error("Exception in thread {}", threadNum, e); - } finally { - endLatch.countDown(); - } - }); - } - - startLatch.countDown(); - assertThat(endLatch.await(10, TimeUnit.SECONDS)).isTrue(); - executor.shutdown(); - executor.awaitTermination(5, TimeUnit.SECONDS); - - if (!exceptions.isEmpty()) { - log.error("=== Member Lock Test Exceptions ==="); - exceptions.forEach(e -> log.error("Exception: {}", e.getMessage(), e)); - } - - assertThat(exceptions).isEmpty(); - assertThat(maxConcurrent.get()).as("최대 동시 실행 수는 1이어야 함").isEqualTo(1); - assertThat(counter.get()).as("카운터는 정확히 " + threadCount + "이어야 함").isEqualTo(threadCount); - } -} diff --git a/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateRepositoryConcurrencyTest.java b/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateRepositoryConcurrencyTest.java index 125ea66d3..94ed344c0 100644 --- a/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateRepositoryConcurrencyTest.java +++ b/src/test/java/com/dongsoop/dongsoop/blinddate/BlindDateRepositoryConcurrencyTest.java @@ -3,7 +3,6 @@ import static org.assertj.core.api.Assertions.assertThat; import com.dongsoop.dongsoop.blinddate.entity.ParticipantInfo; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateSessionLock; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorage; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorageImpl; import com.dongsoop.dongsoop.blinddate.repository.BlindDateSessionStorageImpl; @@ -45,9 +44,8 @@ class BlindDateStorageConcurrencyTest { @BeforeEach void setUp() { - BlindDateSessionLock sessionLock = new BlindDateSessionLock(); participantStorage = new BlindDateParticipantStorageImpl(); - sessionStorage = new BlindDateSessionStorageImpl(sessionLock); + sessionStorage = new BlindDateSessionStorageImpl(); blindDateStorage = new BlindDateStorageImpl(); } diff --git a/src/test/java/com/dongsoop/dongsoop/blinddate/WebSocketTestConfig.java b/src/test/java/com/dongsoop/dongsoop/blinddate/WebSocketTestConfig.java index 19311df91..95d7d48d5 100644 --- a/src/test/java/com/dongsoop/dongsoop/blinddate/WebSocketTestConfig.java +++ b/src/test/java/com/dongsoop/dongsoop/blinddate/WebSocketTestConfig.java @@ -5,9 +5,7 @@ import com.dongsoop.dongsoop.blinddate.handler.BlindDateConnectHandler; import com.dongsoop.dongsoop.blinddate.handler.BlindDateDisconnectHandler; import com.dongsoop.dongsoop.blinddate.handler.BlindDateMessageHandler; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateMatchingLock; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateMemberLock; -import com.dongsoop.dongsoop.blinddate.lock.BlindDateSessionLock; +import com.dongsoop.dongsoop.blinddate.executor.BlindDateEventQueue; import com.dongsoop.dongsoop.blinddate.notification.BlindDateNotification; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorage; import com.dongsoop.dongsoop.blinddate.repository.BlindDateParticipantStorageImpl; @@ -55,20 +53,14 @@ public BlindDateParticipantStorage participantStorage() { @Bean @Primary - public BlindDateSessionStorage sessionStorage(BlindDateSessionLock sessionLock) { - return new BlindDateSessionStorageImpl(sessionLock); + public BlindDateSessionStorage sessionStorage() { + return new BlindDateSessionStorageImpl(); } @Bean @Primary - public BlindDateMatchingLock blindDateMatchingLock() { - return new BlindDateMatchingLock(); - } - - @Bean - @Primary - public BlindDateMemberLock blindDateMemberLock() { - return new BlindDateMemberLock(); + public BlindDateEventQueue blindDateEventQueue() { + return new BlindDateEventQueue(); } @Bean @@ -104,7 +96,8 @@ public BlindDateService blindDateService( BlindDateNotification notification, BlindDateSessionStorage sessionStorage, SimpMessagingTemplate messagingTemplate, - BlindDateTaskScheduler taskScheduler + BlindDateTaskScheduler taskScheduler, + BlindDateEventQueue eventQueue ) { return new BlindDateServiceImpl( participantStorage, @@ -112,7 +105,8 @@ public BlindDateService blindDateService( notification, sessionStorage, messagingTemplate, - taskScheduler + taskScheduler, + eventQueue ); } @@ -138,9 +132,8 @@ public BlindDateConnectHandler blindDateConnectHandler( BlindDateSessionService sessionService, BlindDateSessionScheduler sessionScheduler, SimpMessagingTemplate messagingTemplate, - BlindDateMatchingLock matchingLock, - BlindDateMemberLock memberLock, - BlindDateSessionLock sessionLock + BlindDateEventQueue eventQueue, + BlindDateDisconnectHandler disconnectHandler ) { return new BlindDateConnectHandler( participantStorage, @@ -150,9 +143,8 @@ public BlindDateConnectHandler blindDateConnectHandler( sessionService, sessionScheduler, messagingTemplate, - matchingLock, - memberLock, - sessionLock + eventQueue, + disconnectHandler ); } @@ -162,17 +154,13 @@ public BlindDateDisconnectHandler blindDateDisconnectHandler( BlindDateParticipantStorage participantStorage, BlindDateSessionStorage sessionStorage, BlindDateService blindDateService, - BlindDateMatchingLock matchingLock, - BlindDateMemberLock memberLock, - BlindDateSessionLock sessionLock + BlindDateEventQueue eventQueue ) { return new BlindDateDisconnectHandler( participantStorage, sessionStorage, blindDateService, - matchingLock, - memberLock, - sessionLock + eventQueue ); } diff --git a/src/test/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateLockTest.java b/src/test/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateLockTest.java deleted file mode 100644 index ded4b1c8f..000000000 --- a/src/test/java/com/dongsoop/dongsoop/blinddate/lock/BlindDateLockTest.java +++ /dev/null @@ -1,77 +0,0 @@ -package com.dongsoop.dongsoop.blinddate.lock; -import static org.assertj.core.api.Assertions.assertThat; -import java.util.ArrayList; -import java.util.Collections; -import java.util.List; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicInteger; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.DisplayName; -import org.junit.jupiter.api.Nested; -import org.junit.jupiter.api.RepeatedTest; -import org.junit.jupiter.api.Test; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -@DisplayName("BlindDate Lock 테스트") -class BlindDateLockTest { - private static final Logger log = LoggerFactory.getLogger(BlindDateLockTest.class); - private BlindDateMemberLock memberLock; - private BlindDateSessionLock sessionLock; - @BeforeEach - void setUp() { - memberLock = new BlindDateMemberLock(); - sessionLock = new BlindDateSessionLock(); - } - @Nested - @DisplayName("1. 멤버 Lock 기본 기능") - class MemberLockBasicTests { - @Test - @DisplayName("Lock 획득 및 해제") - void lockAndUnlock() { - Long memberId = 1L; - memberLock.lockByMemberId(memberId); - memberLock.unlockByMemberId(memberId); - } - @Test - @DisplayName("서로 다른 멤버 Lock - 독립적") - void differentMembers_IndependentLocks() throws InterruptedException { - Long member1 = 1L; - Long member2 = 2L; - CountDownLatch latch = new CountDownLatch(2); - AtomicInteger executionOrder = new AtomicInteger(0); - List order = Collections.synchronizedList(new ArrayList<>()); - Thread thread1 = new Thread(() -> { - memberLock.lockByMemberId(member1); - try { - order.add(executionOrder.incrementAndGet()); - Thread.sleep(100); - } catch (InterruptedException e) { - log.error("Thread interrupted", e); - } finally { - memberLock.unlockByMemberId(member1); - latch.countDown(); - } - }); - Thread thread2 = new Thread(() -> { - memberLock.lockByMemberId(member2); - try { - order.add(executionOrder.incrementAndGet()); - Thread.sleep(100); - } catch (InterruptedException e) { - log.error("Thread interrupted", e); - } finally { - memberLock.unlockByMemberId(member2); - latch.countDown(); - } - }); - thread1.start(); - thread2.start(); - assertThat(latch.await(5, TimeUnit.SECONDS)).isTrue(); - assertThat(order).hasSize(2); - } - } -}