diff --git a/src/main/java/gg/agit/konect/domain/chat/repository/ChatRoomMemberRepository.java b/src/main/java/gg/agit/konect/domain/chat/repository/ChatRoomMemberRepository.java index 255030852..9707f1a48 100644 --- a/src/main/java/gg/agit/konect/domain/chat/repository/ChatRoomMemberRepository.java +++ b/src/main/java/gg/agit/konect/domain/chat/repository/ChatRoomMemberRepository.java @@ -79,13 +79,13 @@ List findByChatRoomIdsAndUserId( @Param("userId") Integer userId ); - @Modifying + @Modifying(clearAutomatically = true) @Query(""" UPDATE ChatRoomMember crm SET crm.lastReadAt = :lastReadAt WHERE crm.id.chatRoomId = :chatRoomId AND crm.id.userId = :userId - AND crm.lastReadAt < :lastReadAt + AND (crm.lastReadAt IS NULL OR crm.lastReadAt < :lastReadAt) """) int updateLastReadAtIfOlder( @Param("chatRoomId") Integer chatRoomId, diff --git a/src/main/java/gg/agit/konect/domain/chat/service/ChatRoomMembershipService.java b/src/main/java/gg/agit/konect/domain/chat/service/ChatRoomMembershipService.java index 5e3d6bc30..fb3f3185c 100644 --- a/src/main/java/gg/agit/konect/domain/chat/service/ChatRoomMembershipService.java +++ b/src/main/java/gg/agit/konect/domain/chat/service/ChatRoomMembershipService.java @@ -1,32 +1,53 @@ package gg.agit.konect.domain.chat.service; +import static gg.agit.konect.global.code.ApiResponseCode.FORBIDDEN_CHAT_ROOM_ACCESS; +import static gg.agit.konect.global.code.ApiResponseCode.NOT_FOUND_CHAT_ROOM; + import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.Comparator; +import java.util.List; +import java.util.Map; import java.util.Objects; +import java.util.stream.Collectors; +import org.springframework.dao.DataIntegrityViolationException; +import org.springframework.dao.DuplicateKeyException; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; import org.springframework.transaction.annotation.Transactional; import gg.agit.konect.domain.chat.model.ChatRoom; import gg.agit.konect.domain.chat.model.ChatRoomMember; import gg.agit.konect.domain.chat.repository.ChatRoomMemberRepository; import gg.agit.konect.domain.chat.repository.ChatRoomRepository; +import gg.agit.konect.domain.club.model.Club; import gg.agit.konect.domain.club.model.ClubMember; +import gg.agit.konect.domain.club.repository.ClubMemberRepository; +import gg.agit.konect.domain.user.enums.UserRole; import gg.agit.konect.domain.user.model.User; +import gg.agit.konect.domain.user.repository.UserRepository; +import gg.agit.konect.global.exception.CustomException; import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +@Slf4j @Service @RequiredArgsConstructor @Transactional(readOnly = true) public class ChatRoomMembershipService { + private static final int SYSTEM_ADMIN_ID = 1; + private final ChatRoomRepository chatRoomRepository; private final ChatRoomMemberRepository chatRoomMemberRepository; + private final ClubMemberRepository clubMemberRepository; + private final UserRepository userRepository; @Transactional public void addClubMember(ClubMember clubMember) { LocalDateTime baseline = Objects.requireNonNull(clubMember.getCreatedAt()); - ChatRoom room = chatRoomRepository.findByClubId(clubMember.getClub().getId()) - .orElseGet(() -> chatRoomRepository.save(ChatRoom.groupOf(clubMember.getClub()))); + ChatRoom room = findOrCreateClubRoom(clubMember.getClub()); ensureMember(room, clubMember.getUser(), baseline); } @@ -37,6 +58,147 @@ public void addDirectMembers(ChatRoom room, User firstUser, User secondUser, Loc ensureMember(room, secondUser, baseline); } + @Transactional + public void removeClubMember(Integer clubId, Integer userId) { + chatRoomRepository.findByClubId(clubId) + .ifPresent(room -> chatRoomMemberRepository.deleteByChatRoomIdAndUserId(room.getId(), userId)); + } + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void ensureClubRoomMemberships(Integer userId) { + List memberships = clubMemberRepository.findAllByUserId(userId); + if (memberships.isEmpty()) { + return; + } + + Map membershipByClubId = memberships.stream() + .collect(Collectors.toMap(cm -> cm.getClub().getId(), cm -> cm, (a, b) -> a)); + + List rooms = resolveOrCreateClubRooms(memberships).stream() + .sorted(Comparator.comparing(ChatRoom::getId)) + .toList(); + List roomIds = rooms.stream().map(ChatRoom::getId).toList(); + if (roomIds.isEmpty()) { + return; + } + + Map memberByRoomId = chatRoomMemberRepository + .findByChatRoomIdsAndUserId(roomIds, userId) + .stream() + .collect(Collectors.toMap(ChatRoomMember::getChatRoomId, member -> member, (a, b) -> a)); + + for (ChatRoom room : rooms) { + ClubMember member = membershipByClubId.get(room.getClub().getId()); + if (member == null) { + continue; + } + + ChatRoomMember existingMember = memberByRoomId.get(room.getId()); + if (existingMember != null) { + LocalDateTime lastReadAt = existingMember.getLastReadAt(); + if (lastReadAt == null || lastReadAt.isBefore(member.getCreatedAt())) { + chatRoomMemberRepository.updateLastReadAtIfOlder( + room.getId(), userId, member.getCreatedAt() + ); + } + continue; + } + + saveRoomMemberIgnoringDuplicate(room, member.getUser(), member.getCreatedAt()); + } + } + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void updateLastReadAt(Integer roomId, Integer userId, LocalDateTime readAt) { + chatRoomMemberRepository.updateLastReadAtIfOlder(roomId, userId, readAt); + } + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void updateDirectRoomLastReadAt(Integer roomId, Integer userId, LocalDateTime readAt) { + User user = userRepository.getById(userId); + ChatRoom room = chatRoomRepository.findById(roomId) + .orElseThrow(() -> CustomException.of(NOT_FOUND_CHAT_ROOM)); + + ensureDirectRoomMemberExists(room, user, readAt); + + if (user.getRole() == UserRole.ADMIN) { + List members = chatRoomMemberRepository.findByChatRoomId(roomId); + boolean isSystemAdmin = members.stream() + .anyMatch(member -> Objects.equals(member.getUserId(), SYSTEM_ADMIN_ID)); + + if (isSystemAdmin) { + for (ChatRoomMember member : members) { + if (member.getUser().getRole() == UserRole.ADMIN) { + chatRoomMemberRepository.updateLastReadAtIfOlder(roomId, member.getUserId(), readAt); + } + } + return; + } + } + + chatRoomMemberRepository.updateLastReadAtIfOlder(roomId, userId, readAt); + } + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void ensureClubRoomMember(Integer roomId, Integer userId) { + ChatRoom room = chatRoomRepository.findById(roomId) + .orElseThrow(() -> CustomException.of(NOT_FOUND_CHAT_ROOM)); + if (!room.isGroupRoom() || room.getClub() == null) { + throw CustomException.of(NOT_FOUND_CHAT_ROOM); + } + ClubMember member = clubMemberRepository.getByClubIdAndUserId(room.getClub().getId(), userId); + ensureMember(room, member.getUser(), member.getCreatedAt()); + } + + private ChatRoom findOrCreateClubRoom(Club club) { + return chatRoomRepository.findByClubId(club.getId()) + .orElseGet(() -> { + try { + return chatRoomRepository.save(ChatRoom.groupOf(club)); + } catch (DataIntegrityViolationException e) { + if (!isDuplicateKeyException(e)) { + throw e; + } + log.debug("클럽 채팅방 동시 생성 감지, 재조회: clubId={}", club.getId()); + return chatRoomRepository.findByClubId(club.getId()) + .orElseThrow(() -> CustomException.of(NOT_FOUND_CHAT_ROOM)); + } + }); + } + + private List resolveOrCreateClubRooms(List memberships) { + Map clubById = memberships.stream() + .map(ClubMember::getClub) + .collect(Collectors.toMap(Club::getId, club -> club, (a, b) -> a)); + + Map roomByClubId = chatRoomRepository.findByClubIds(new ArrayList<>(clubById.keySet())) + .stream() + .filter(room -> room.getClub() != null) + .collect(Collectors.toMap(room -> room.getClub().getId(), room -> room, (a, b) -> a)); + + for (Map.Entry clubEntry : clubById.entrySet()) { + if (roomByClubId.containsKey(clubEntry.getKey())) { + continue; + } + try { + ChatRoom createdRoom = chatRoomRepository.save(ChatRoom.groupOf(clubEntry.getValue())); + roomByClubId.put(clubEntry.getKey(), createdRoom); + } catch (DataIntegrityViolationException e) { + if (!isDuplicateKeyException(e)) { + throw e; + } + log.debug("클럽 채팅방 동시 생성 감지, 재조회: clubId={}", clubEntry.getKey()); + chatRoomRepository.findByClubId(clubEntry.getKey()) + .ifPresent(room -> roomByClubId.put(clubEntry.getKey(), room)); + } + } + + return memberships.stream() + .map(membership -> roomByClubId.get(membership.getClub().getId())) + .filter(Objects::nonNull) + .toList(); + } + private void ensureMember(ChatRoom room, User user, LocalDateTime baseline) { chatRoomMemberRepository.findByChatRoomIdAndUserId(room.getId(), user.getId()) .ifPresentOrElse(member -> { @@ -44,12 +206,50 @@ private void ensureMember(ChatRoom room, User user, LocalDateTime baseline) { if (lastReadAt == null || lastReadAt.isBefore(baseline)) { member.updateLastReadAt(baseline); } - }, () -> chatRoomMemberRepository.save(ChatRoomMember.of(room, user, baseline))); + }, () -> saveRoomMemberIgnoringDuplicate(room, user, baseline)); } - @Transactional - public void removeClubMember(Integer clubId, Integer userId) { - chatRoomRepository.findByClubId(clubId) - .ifPresent(room -> chatRoomMemberRepository.deleteByChatRoomIdAndUserId(room.getId(), userId)); + private void saveRoomMemberIgnoringDuplicate(ChatRoom room, User user, LocalDateTime baseline) { + try { + chatRoomMemberRepository.save(ChatRoomMember.of(room, user, baseline)); + } catch (DataIntegrityViolationException e) { + if (!isDuplicateKeyException(e)) { + throw e; + } + log.debug("채팅방 멤버 동시 생성 감지, 무시: roomId={}, userId={}", room.getId(), user.getId()); + } + } + + private void ensureDirectRoomMemberExists(ChatRoom room, User user, LocalDateTime readAt) { + boolean exists = chatRoomMemberRepository.existsByChatRoomIdAndUserId(room.getId(), user.getId()); + if (exists) { + return; + } + + if (user.getRole() == UserRole.ADMIN && isSystemAdminRoom(room.getId())) { + saveRoomMemberIgnoringDuplicate(room, user, readAt); + return; + } + + throw CustomException.of(FORBIDDEN_CHAT_ROOM_ACCESS); + } + + private boolean isSystemAdminRoom(Integer roomId) { + List memberIds = chatRoomMemberRepository.findRoomMemberIdsByChatRoomIds(List.of(roomId)); + return memberIds.stream() + .map(row -> (Integer)row[1]) + .anyMatch(userId -> userId.equals(SYSTEM_ADMIN_ID)); + } + + private boolean isDuplicateKeyException(DataIntegrityViolationException e) { + if (e instanceof DuplicateKeyException) { + return true; + } + Throwable rootCause = e.getRootCause(); + if (rootCause == null) { + return false; + } + String message = rootCause.getMessage(); + return message != null && (message.contains("Duplicate") || message.contains("duplicate key")); } } diff --git a/src/main/java/gg/agit/konect/domain/chat/service/ChatService.java b/src/main/java/gg/agit/konect/domain/chat/service/ChatService.java index da2c27320..0d5cc2c24 100644 --- a/src/main/java/gg/agit/konect/domain/chat/service/ChatService.java +++ b/src/main/java/gg/agit/konect/domain/chat/service/ChatService.java @@ -20,6 +20,7 @@ import org.springframework.data.domain.Page; import org.springframework.data.domain.PageRequest; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; import org.springframework.transaction.annotation.Transactional; import org.springframework.util.StringUtils; @@ -43,7 +44,6 @@ import gg.agit.konect.domain.chat.repository.ChatRoomMemberRepository; import gg.agit.konect.domain.chat.repository.ChatRoomRepository; import gg.agit.konect.domain.chat.repository.RoomUnreadCountProjection; -import gg.agit.konect.domain.club.model.Club; import gg.agit.konect.domain.club.model.ClubMember; import gg.agit.konect.domain.club.repository.ClubMemberRepository; import gg.agit.konect.domain.notification.enums.NotificationTargetType; @@ -56,7 +56,9 @@ import gg.agit.konect.global.code.ApiResponseCode; import gg.agit.konect.global.exception.CustomException; import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +@Slf4j @Service @RequiredArgsConstructor @Transactional(readOnly = true) @@ -71,6 +73,7 @@ public class ChatService { private final ClubMemberRepository clubMemberRepository; private final UserRepository userRepository; private final ChatPresenceService chatPresenceService; + private final ChatRoomMembershipService chatRoomMembershipService; private final NotificationService notificationService; private final ApplicationEventPublisher eventPublisher; @@ -148,8 +151,9 @@ public void leaveChatRoom(Integer userId, Integer roomId) { chatRoomMemberRepository.deleteByChatRoomIdAndUserId(roomId, userId); } - @Transactional public ChatRoomsSummaryResponse getChatRooms(Integer userId) { + chatRoomMembershipService.ensureClubRoomMemberships(userId); + List directRooms = getDirectChatRooms(userId); List clubRooms = getClubChatRooms(userId); @@ -191,15 +195,22 @@ public ChatRoomsSummaryResponse getChatRooms(Integer userId) { return new ChatRoomsSummaryResponse(rooms); } - @Transactional + @Transactional(propagation = Propagation.NOT_SUPPORTED) public ChatMessagePageResponse getMessages(Integer userId, Integer roomId, Integer page, Integer limit) { ChatRoom room = chatRoomRepository.findById(roomId) .orElseThrow(() -> CustomException.of(NOT_FOUND_CHAT_ROOM)); + LocalDateTime readAt = LocalDateTime.now(); + if (room.isDirectRoom()) { - return getDirectChatRoomMessages(userId, roomId, page, limit); + chatRoomMembershipService.updateDirectRoomLastReadAt(roomId, userId, readAt); + recordPresenceSafely(roomId, userId); + return getDirectChatRoomMessages(userId, roomId, page, limit, readAt); } + chatRoomMembershipService.ensureClubRoomMember(roomId, userId); + chatRoomMembershipService.updateLastReadAt(roomId, userId, readAt); + recordPresenceSafely(roomId, userId); return getClubMessagesByRoomId(roomId, userId, page, limit); } @@ -324,32 +335,60 @@ private List getAdminDirectChatRooms() { .toList(); } + private List getClubChatRooms(Integer userId) { + List memberships = clubMemberRepository.findAllByUserId(userId); + if (memberships.isEmpty()) { + return List.of(); + } + + List clubIds = memberships.stream() + .map(cm -> cm.getClub().getId()) + .toList(); + + List rooms = chatRoomRepository.findByClubIds(new ArrayList<>(clubIds)) + .stream() + .filter(room -> room.getClub() != null) + .toList(); + + List roomIds = rooms.stream().map(ChatRoom::getId).toList(); + Map lastMessageMap = getLastMessageMap(roomIds); + Map unreadCountMap = getRoomUnreadCountMap(roomIds, userId); + + return rooms.stream() + .map(room -> { + ChatMessage lastMessage = lastMessageMap.get(room.getId()); + return new ChatRoomSummaryResponse( + room.getId(), + ChatType.GROUP, + room.getClub().getName(), + room.getClub().getImageUrl(), + lastMessage != null ? lastMessage.getContent() : null, + lastMessage != null ? lastMessage.getCreatedAt() : null, + unreadCountMap.getOrDefault(room.getId(), 0), + false + ); + }) + .toList(); + } + private ChatMessagePageResponse getDirectChatRoomMessages( Integer userId, Integer roomId, Integer page, - Integer limit + Integer limit, + LocalDateTime readAt ) { ChatRoom chatRoom = getDirectRoom(roomId); User user = userRepository.getById(userId); ChatRoomMember member = getOrCreateDirectRoomMember(chatRoom, user); LocalDateTime visibleMessageFrom = prepareDirectRoomAccess(member, chatRoom); - LocalDateTime readAt = LocalDateTime.now(); - chatPresenceService.recordPresence(roomId, userId); - boolean isAdminViewingSystemRoom = user.getRole() == UserRole.ADMIN && isSystemAdminRoom(chatRoom); PageRequest pageable = PageRequest.of(page - 1, limit); Page messages = chatMessageRepository.findByChatRoomId(roomId, visibleMessageFrom, pageable); List members = chatRoomMemberRepository.findByChatRoomId(roomId); - if (isAdminViewingSystemRoom) { - updateAllAdminMembersLastReadAt(members, readAt); - } else { - member.updateLastReadAt(readAt); - } - List sortedReadBaselines = isAdminViewingSystemRoom ? toAdminChatReadBaselines(members) : toSortedReadBaselines(members); @@ -394,7 +433,7 @@ private ChatMessageDetailResponse sendDirectMessage( ChatRoomMember senderMember = getAccessibleDirectRoomMember(chatRoom, sender); boolean senderHadLeft = senderMember.hasLeft(); List members = chatRoomMemberRepository.findByChatRoomId(roomId); - User receiver = resolveMessageReceiver(sender, members); + User receiver = resolveDirectChatPartner(members, userId); ChatMessage chatMessage = chatMessageRepository.save( ChatMessage.of(chatRoom, sender, request.content()) @@ -425,39 +464,6 @@ private ChatMessageDetailResponse sendDirectMessage( ); } - private List getClubChatRooms(Integer userId) { - List memberships = clubMemberRepository.findAllByUserId(userId); - if (memberships.isEmpty()) { - return List.of(); - } - - Map membershipByClubId = memberships.stream() - .collect(Collectors.toMap(cm -> cm.getClub().getId(), cm -> cm, (a, b) -> a)); - - List rooms = resolveOrCreateClubRooms(memberships); - ensureClubRoomMembers(rooms, membershipByClubId, userId); - - List roomIds = rooms.stream().map(ChatRoom::getId).toList(); - Map lastMessageMap = getLastMessageMap(roomIds); - Map unreadCountMap = getRoomUnreadCountMap(roomIds, userId); - - return rooms.stream() - .map(room -> { - ChatMessage lastMessage = lastMessageMap.get(room.getId()); - return new ChatRoomSummaryResponse( - room.getId(), - ChatType.GROUP, - room.getClub().getName(), - room.getClub().getImageUrl(), - lastMessage != null ? lastMessage.getContent() : null, - lastMessage != null ? lastMessage.getCreatedAt() : null, - unreadCountMap.getOrDefault(room.getId(), 0), - false - ); - }) - .toList(); - } - private ChatMessagePageResponse getClubMessagesByRoomId( Integer roomId, Integer userId, @@ -465,11 +471,6 @@ private ChatMessagePageResponse getClubMessagesByRoomId( Integer limit ) { ChatRoom room = getClubRoom(roomId); - ClubMember member = clubMemberRepository.getByClubIdAndUserId(room.getClub().getId(), userId); - ensureRoomMember(room, member.getUser(), member.getCreatedAt()); - - chatPresenceService.recordPresence(roomId, userId); - updateLastReadAt(roomId, userId, LocalDateTime.now()); PageRequest pageable = PageRequest.of(page - 1, limit); long totalCount = chatMessageRepository.countByChatRoomId(roomId, null); @@ -594,63 +595,6 @@ private ChatRoom getClubRoom(Integer roomId) { return room; } - private List resolveOrCreateClubRooms(List memberships) { - Map clubById = memberships.stream() - .map(ClubMember::getClub) - .collect(Collectors.toMap(Club::getId, club -> club, (a, b) -> a)); - - Map roomByClubId = chatRoomRepository.findByClubIds(new ArrayList<>(clubById.keySet())) - .stream() - .filter(room -> room.getClub() != null) - .collect(Collectors.toMap(room -> room.getClub().getId(), room -> room, (a, b) -> a)); - - for (Map.Entry clubEntry : clubById.entrySet()) { - if (roomByClubId.containsKey(clubEntry.getKey())) { - continue; - } - - ChatRoom createdRoom = chatRoomRepository.save(ChatRoom.groupOf(clubEntry.getValue())); - roomByClubId.put(clubEntry.getKey(), createdRoom); - } - - return memberships.stream() - .map(membership -> roomByClubId.get(membership.getClub().getId())) - .toList(); - } - - private void ensureClubRoomMembers( - List rooms, - Map membershipByClubId, - Integer userId - ) { - if (rooms.isEmpty()) { - return; - } - - Map memberByRoomId = chatRoomMemberRepository - .findByChatRoomIdsAndUserId(extractChatRoomIds(rooms), userId) - .stream() - .collect(Collectors.toMap(ChatRoomMember::getChatRoomId, member -> member, (a, b) -> a)); - - for (ChatRoom room : rooms) { - ClubMember member = membershipByClubId.get(room.getClub().getId()); - if (member == null) { - continue; - } - - ChatRoomMember existingMember = memberByRoomId.get(room.getId()); - if (existingMember != null) { - LocalDateTime lastReadAt = existingMember.getLastReadAt(); - if (lastReadAt == null || lastReadAt.isBefore(member.getCreatedAt())) { - existingMember.updateLastReadAt(member.getCreatedAt()); - } - continue; - } - - chatRoomMemberRepository.save(ChatRoomMember.of(room, member.getUser(), member.getCreatedAt())); - } - } - private List extractChatRoomIds(List chatRooms) { return chatRooms.stream() .map(ChatRoom::getId) @@ -674,23 +618,6 @@ private Map getUnreadCountMap(List chatRoomIds, Integ )); } - private Map getAdminUnreadCountMap(List chatRoomIds) { - if (chatRoomIds.isEmpty()) { - return Map.of(); - } - - List unreadMessageCounts = chatMessageRepository.countUnreadMessagesForAdmin( - chatRoomIds, - UserRole.ADMIN - ); - - return unreadMessageCounts.stream() - .collect(Collectors.toMap( - UnreadMessageCount::chatRoomId, - unreadMessageCount -> unreadMessageCount.unreadCount().intValue() - )); - } - private Integer getMaskedAdminId(User user, ChatRoom chatRoom) { if (user.getRole() == UserRole.ADMIN) { return null; @@ -820,14 +747,6 @@ private List toAdminChatReadBaselines(List member return baselines; } - private void updateAllAdminMembersLastReadAt(List members, LocalDateTime readAt) { - for (ChatRoomMember member : members) { - if (member.getUser().getRole() == UserRole.ADMIN) { - member.updateLastReadAt(readAt); - } - } - } - private int countUnreadSince(LocalDateTime messageCreatedAt, List sortedReadBaselines) { int left = 0; int right = sortedReadBaselines.size(); @@ -1055,14 +974,6 @@ private User resolveDirectChatPartner( return findDirectPartnerFromMemberInfo(memberInfos, userId, userMap); } - private User findNonAdminMember(List members) { - return members.stream() - .map(ChatRoomMember::getUser) - .filter(memberUser -> memberUser.getRole() != UserRole.ADMIN) - .findFirst() - .orElse(null); - } - private User findNonAdminUserFromMemberInfo(List memberInfos, Map userMap) { return memberInfos.stream() .sorted(Comparator.comparing(MemberInfo::createdAt)) @@ -1073,21 +984,6 @@ private User findNonAdminUserFromMemberInfo(List memberInfos, Map members) { - if (sender.getRole() == UserRole.ADMIN) { - User nonAdminUser = findNonAdminMember(members); - if (nonAdminUser != null) { - return nonAdminUser; - } - } - - User partner = findDirectPartner(members, sender.getId()); - if (partner == null) { - throw CustomException.of(FORBIDDEN_CHAT_ROOM_ACCESS); - } - return partner; - } - private User resolveMessageReceiverFromMemberInfo( User sender, List memberInfos, @@ -1107,4 +1003,11 @@ private User resolveMessageReceiverFromMemberInfo( return partner; } + private void recordPresenceSafely(Integer roomId, Integer userId) { + try { + chatPresenceService.recordPresence(roomId, userId); + } catch (Exception e) { + log.warn("Redis presence record failed, continuing: roomId={}, userId={}", roomId, userId, e); + } + } } diff --git a/src/main/java/gg/agit/konect/domain/notification/service/NotificationService.java b/src/main/java/gg/agit/konect/domain/notification/service/NotificationService.java index a597771e6..c09cdf499 100644 --- a/src/main/java/gg/agit/konect/domain/notification/service/NotificationService.java +++ b/src/main/java/gg/agit/konect/domain/notification/service/NotificationService.java @@ -73,7 +73,7 @@ public void deleteToken(Integer userId, NotificationTokenDeleteRequest request) .ifPresent(notificationDeviceTokenRepository::delete); } - @Async + @Async("notificationTaskExecutor") @Transactional public void sendChatNotification(Integer receiverId, Integer roomId, String senderName, String messageContent) { try { @@ -123,7 +123,7 @@ public void sendChatNotification(Integer receiverId, Integer roomId, String send } } - @Async + @Async("notificationTaskExecutor") @Transactional public void sendGroupChatNotification( Integer roomId, @@ -210,7 +210,7 @@ public void sendGroupChatNotification( } } - @Async + @Async("notificationTaskExecutor") @Transactional public void sendClubApplicationSubmittedNotification( Integer receiverId, @@ -227,7 +227,7 @@ public void sendClubApplicationSubmittedNotification( sendNotification(receiverId, clubName, body, path); } - @Async + @Async("notificationTaskExecutor") @Transactional public void sendClubApplicationApprovedNotification(Integer receiverId, Integer clubId, String clubName) { String body = "동아리 지원이 승인되었어요."; @@ -238,7 +238,7 @@ public void sendClubApplicationApprovedNotification(Integer receiverId, Integer sendNotification(receiverId, clubName, body, path); } - @Async + @Async("notificationTaskExecutor") @Transactional public void sendClubApplicationRejectedNotification(Integer receiverId, Integer clubId, String clubName) { String body = "동아리 지원이 거절되었어요."; diff --git a/src/main/java/gg/agit/konect/global/config/AsyncConfig.java b/src/main/java/gg/agit/konect/global/config/AsyncConfig.java index bfc63895d..f88deb3ef 100644 --- a/src/main/java/gg/agit/konect/global/config/AsyncConfig.java +++ b/src/main/java/gg/agit/konect/global/config/AsyncConfig.java @@ -1,20 +1,61 @@ package gg.agit.konect.global.config; import java.util.concurrent.Executor; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ThreadPoolExecutor; +import org.springframework.aop.interceptor.AsyncUncaughtExceptionHandler; +import org.springframework.aop.interceptor.SimpleAsyncUncaughtExceptionHandler; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.annotation.AsyncConfigurer; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import lombok.extern.slf4j.Slf4j; + +@Slf4j @Configuration -public class AsyncConfig { +public class AsyncConfig implements AsyncConfigurer { + + private static final int DEFAULT_CORE_POOL_SIZE = 2; + private static final int DEFAULT_MAX_POOL_SIZE = 5; + private static final int DEFAULT_QUEUE_CAPACITY = 50; + private static final int DEFAULT_AWAIT_TERMINATION_SECONDS = 30; private static final int SHEET_SYNC_CORE_POOL_SIZE = 2; private static final int SHEET_SYNC_MAX_POOL_SIZE = 4; private static final int SHEET_SYNC_QUEUE_CAPACITY = 50; private static final int SHEET_SYNC_AWAIT_TERMINATION_SECONDS = 30; + private static final int NOTIFICATION_CORE_POOL_SIZE = 2; + private static final int NOTIFICATION_MAX_POOL_SIZE = 5; + private static final int NOTIFICATION_QUEUE_CAPACITY = 100; + private static final int NOTIFICATION_AWAIT_TERMINATION_SECONDS = 30; + + private static final int SLACK_CORE_POOL_SIZE = 1; + private static final int SLACK_MAX_POOL_SIZE = 3; + private static final int SLACK_QUEUE_CAPACITY = 50; + private static final int SLACK_AWAIT_TERMINATION_SECONDS = 30; + + @Override + public Executor getAsyncExecutor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(DEFAULT_CORE_POOL_SIZE); + executor.setMaxPoolSize(DEFAULT_MAX_POOL_SIZE); + executor.setQueueCapacity(DEFAULT_QUEUE_CAPACITY); + executor.setThreadNamePrefix("async-default-"); + executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); + executor.setWaitForTasksToCompleteOnShutdown(true); + executor.setAwaitTerminationSeconds(DEFAULT_AWAIT_TERMINATION_SECONDS); + executor.initialize(); + return executor; + } + + @Override + public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() { + return new SimpleAsyncUncaughtExceptionHandler(); + } + @Bean(name = "sheetSyncTaskExecutor") public Executor sheetSyncTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); @@ -28,4 +69,40 @@ public Executor sheetSyncTaskExecutor() { executor.initialize(); return executor; } + + @Bean(name = "notificationTaskExecutor") + public Executor notificationTaskExecutor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(NOTIFICATION_CORE_POOL_SIZE); + executor.setMaxPoolSize(NOTIFICATION_MAX_POOL_SIZE); + executor.setQueueCapacity(NOTIFICATION_QUEUE_CAPACITY); + executor.setThreadNamePrefix("notification-"); + executor.setRejectedExecutionHandler((runnable, pool) -> { + log.warn("알림 스레드풀 포화로 작업이 거절되었습니다. poolSize={}, activeCount={}, queueSize={}", + pool.getPoolSize(), pool.getActiveCount(), pool.getQueue().size()); + throw new RejectedExecutionException("notificationTaskExecutor saturated"); + }); + executor.setWaitForTasksToCompleteOnShutdown(true); + executor.setAwaitTerminationSeconds(NOTIFICATION_AWAIT_TERMINATION_SECONDS); + executor.initialize(); + return executor; + } + + @Bean(name = "slackTaskExecutor") + public Executor slackTaskExecutor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(SLACK_CORE_POOL_SIZE); + executor.setMaxPoolSize(SLACK_MAX_POOL_SIZE); + executor.setQueueCapacity(SLACK_QUEUE_CAPACITY); + executor.setThreadNamePrefix("slack-"); + executor.setRejectedExecutionHandler((runnable, pool) -> { + log.warn("Slack 스레드풀 포화로 작업이 거절되었습니다. poolSize={}, activeCount={}, queueSize={}", + pool.getPoolSize(), pool.getActiveCount(), pool.getQueue().size()); + throw new RejectedExecutionException("slackTaskExecutor saturated"); + }); + executor.setWaitForTasksToCompleteOnShutdown(true); + executor.setAwaitTerminationSeconds(SLACK_AWAIT_TERMINATION_SECONDS); + executor.initialize(); + return executor; + } } diff --git a/src/main/java/gg/agit/konect/infrastructure/slack/ai/SlackAIService.java b/src/main/java/gg/agit/konect/infrastructure/slack/ai/SlackAIService.java index c984ef59f..d050a21e9 100644 --- a/src/main/java/gg/agit/konect/infrastructure/slack/ai/SlackAIService.java +++ b/src/main/java/gg/agit/konect/infrastructure/slack/ai/SlackAIService.java @@ -81,7 +81,7 @@ public List> fetchAIThreadReplies(String channelId, String t return new ArrayList<>(); } - @Async + @Async("slackTaskExecutor") public void processAIQuery(String text, String channelId, String threadTs, List> cachedReplies) { try { diff --git a/src/main/java/gg/agit/konect/infrastructure/slack/listener/ChatSlackListener.java b/src/main/java/gg/agit/konect/infrastructure/slack/listener/ChatSlackListener.java index c2b7d9aef..44b6774ce 100644 --- a/src/main/java/gg/agit/konect/infrastructure/slack/listener/ChatSlackListener.java +++ b/src/main/java/gg/agit/konect/infrastructure/slack/listener/ChatSlackListener.java @@ -16,7 +16,7 @@ public class ChatSlackListener { private final SlackNotificationService slackNotificationService; - @Async + @Async("slackTaskExecutor") @TransactionalEventListener(phase = AFTER_COMMIT) public void handleAdminChatReceived(AdminChatReceivedEvent event) { slackNotificationService.notifyAdminChatReceived(event.senderName(), event.content()); diff --git a/src/main/java/gg/agit/konect/infrastructure/slack/listener/InquirySlackListener.java b/src/main/java/gg/agit/konect/infrastructure/slack/listener/InquirySlackListener.java index 578c2b22e..32f98195e 100644 --- a/src/main/java/gg/agit/konect/infrastructure/slack/listener/InquirySlackListener.java +++ b/src/main/java/gg/agit/konect/infrastructure/slack/listener/InquirySlackListener.java @@ -16,7 +16,7 @@ public class InquirySlackListener { private final SlackNotificationService slackNotificationService; - @Async + @Async("slackTaskExecutor") @TransactionalEventListener(phase = AFTER_COMMIT) public void handleInquirySubmitted(InquirySubmittedEvent event) { slackNotificationService.notifyInquiry(event.content()); diff --git a/src/main/java/gg/agit/konect/infrastructure/slack/listener/UserSlackListener.java b/src/main/java/gg/agit/konect/infrastructure/slack/listener/UserSlackListener.java index eea11dc88..8b5caa244 100644 --- a/src/main/java/gg/agit/konect/infrastructure/slack/listener/UserSlackListener.java +++ b/src/main/java/gg/agit/konect/infrastructure/slack/listener/UserSlackListener.java @@ -17,13 +17,13 @@ public class UserSlackListener { private final SlackNotificationService slackNotificationService; - @Async + @Async("slackTaskExecutor") @TransactionalEventListener(phase = AFTER_COMMIT) public void handleUserWithdrawn(UserWithdrawnEvent event) { slackNotificationService.notifyUserWithdraw(event.email(), event.provider()); } - @Async + @Async("slackTaskExecutor") @TransactionalEventListener(phase = AFTER_COMMIT) public void handleUserRegistered(UserRegisteredEvent event) { slackNotificationService.notifyUserRegister(event.email(), event.provider()); diff --git a/src/main/resources/application-db.yml b/src/main/resources/application-db.yml index 404202a1f..de2bb656e 100644 --- a/src/main/resources/application-db.yml +++ b/src/main/resources/application-db.yml @@ -9,6 +9,11 @@ spring: url: ${MYSQL_URL} username: ${MYSQL_USERNAME} password: ${MYSQL_PASSWORD} + hikari: + maximum-pool-size: 20 + minimum-idle: 10 + connection-timeout: 5000 + leak-detection-threshold: 30000 jpa: properties: