diff --git a/backend/src/main/java/com/opensource/docgrid/domain/auth/service/query/PrimaryRoleQueryService.java b/backend/src/main/java/com/opensource/docgrid/domain/auth/service/query/PrimaryRoleQueryService.java
index 3d72eab6..c40dac9d 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/auth/service/query/PrimaryRoleQueryService.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/auth/service/query/PrimaryRoleQueryService.java
@@ -1,6 +1,8 @@
package com.opensource.docgrid.domain.auth.service.query;
+import java.util.HashSet;
import java.util.List;
+import java.util.Set;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Propagation;
@@ -12,7 +14,7 @@
import lombok.RequiredArgsConstructor;
/**
- * HTTP 관리자 인가에 필요한 역할을 Redis나 standby 없이 현재 primary에서 확인한다.
+ * HTTP 관리자 인가와 대시보드 push에 필요한 역할을 Redis나 standby 없이 현재 primary에서 확인한다.
*
*
호출자의 read-only 트랜잭션을 상속하지 않으며, OpenProxy가 잘못 라우팅하면
* 권한을 추정하지 않고 요청을 실패시킨다. 일반 API·WebSocket의 역할 캐시는 담당하지 않는다.
@@ -21,19 +23,44 @@
@RequiredArgsConstructor
public class PrimaryRoleQueryService {
+ private static final int MAX_BATCH_SIZE = 500;
+
private final EntityManager entityManager;
private final UserRoleRepository userRoleRepository;
@Transactional(propagation = Propagation.REQUIRES_NEW)
public List findCurrentRoles(Long userId) {
// 1. 명시적 read-write 트랜잭션 안에서 OpenProxy의 실제 도착 역할을 확인한다.
+ requirePrimary();
+
+ // 2. 같은 트랜잭션의 최신 primary에서 역할을 조회한다.
+ return userRoleRepository.findRoleCodesByUserId(userId);
+ }
+
+ @Transactional(propagation = Propagation.REQUIRES_NEW)
+ public Set findCurrentAdminUserIds(List userIds) {
+ if (userIds.isEmpty()) {
+ return Set.of();
+ }
+
+ // 1. 단일 read-write 트랜잭션에서 primary를 확인해 standby의 낡은 ADMIN을 쓰지 않는다.
+ requirePrimary();
+
+ // 2. 큰 세션 집합도 크기가 제한된 IN 쿼리로 읽되, 결과는 이 push에만 쓰는 불변 집합으로 만든다.
+ Set admins = new HashSet<>();
+ for (int start = 0; start < userIds.size(); start += MAX_BATCH_SIZE) {
+ admins.addAll(userRoleRepository.findAdminUserIdsByUserIdIn(
+ userIds.subList(start, Math.min(start + MAX_BATCH_SIZE, userIds.size()))
+ ));
+ }
+ return Set.copyOf(admins);
+ }
+
+ private void requirePrimary() {
Object inRecovery = entityManager.createNativeQuery("SELECT pg_is_in_recovery()")
.getSingleResult();
if (!Boolean.FALSE.equals(inRecovery)) {
throw new IllegalStateException("Primary role verification is unavailable");
}
-
- // 2. 같은 트랜잭션의 최신 primary에서 역할을 조회한다.
- return userRoleRepository.findRoleCodesByUserId(userId);
}
}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/auth/websocket/DashboardAuthorizationSnapshot.java b/backend/src/main/java/com/opensource/docgrid/domain/auth/websocket/DashboardAuthorizationSnapshot.java
new file mode 100644
index 00000000..e045fc81
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/auth/websocket/DashboardAuthorizationSnapshot.java
@@ -0,0 +1,36 @@
+package com.opensource.docgrid.domain.auth.websocket;
+
+import java.util.Map;
+import java.util.Set;
+
+/**
+ * 대시보드 push 한 번에 대해 후보 물리 세션과 primary에서 확인한 ADMIN 사용자만 내부 메시지에 전달한다.
+ *
+ *
브로커가 복제한 수신자별 MESSAGE에서만 읽으며 클라이언트 native STOMP 헤더에는 싣지 않는다.
+ * 다른 push에 재사용하지 않고, 로그에 사용자 ID가 노출되지 않도록 문자열 표현도 제한한다.
+ */
+public record DashboardAuthorizationSnapshot(Map candidateSessions, Set adminUserIds) {
+
+ public static final String HEADER = "docgrid.dashboard.authorizationSnapshot";
+
+ public DashboardAuthorizationSnapshot {
+ candidateSessions = Map.copyOf(candidateSessions);
+ adminUserIds = Set.copyOf(adminUserIds);
+ if (!Set.copyOf(candidateSessions.values()).containsAll(adminUserIds)) {
+ throw new IllegalArgumentException("ADMIN 판정은 이번 push의 후보 사용자에 한정해야 합니다.");
+ }
+ }
+
+ public boolean wasCandidate(String sessionId, Long userId) {
+ return userId.equals(candidateSessions.get(sessionId));
+ }
+
+ public boolean allows(String sessionId, Long userId) {
+ return wasCandidate(sessionId, userId) && adminUserIds.contains(userId);
+ }
+
+ @Override
+ public String toString() {
+ return "DashboardAuthorizationSnapshot[redacted]";
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/auth/websocket/StompDashboardOutboundAuthorizationInterceptor.java b/backend/src/main/java/com/opensource/docgrid/domain/auth/websocket/StompDashboardOutboundAuthorizationInterceptor.java
index 0bf20aab..e9605ee1 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/auth/websocket/StompDashboardOutboundAuthorizationInterceptor.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/auth/websocket/StompDashboardOutboundAuthorizationInterceptor.java
@@ -7,27 +7,24 @@
import org.springframework.messaging.support.ChannelInterceptor;
import org.springframework.stereotype.Component;
-import com.opensource.docgrid.domain.auth.jwt.RoleAuthorityService;
-
import lombok.RequiredArgsConstructor;
-import lombok.extern.slf4j.Slf4j;
/**
- * 기존 대시보드 구독에 대한 실제 outbound MESSAGE 전송을 현재 primary 역할로 제한한다.
+ * 기존 대시보드 구독에 대한 실제 outbound MESSAGE 전송을 push별 primary 판정으로 제한한다.
*
*
SUBSCRIBE 이후 역할이 회수되어도 브로커에는 구독이 남을 수 있다. 이 경계는 각 물리 세션에
- * 전달할 때 다시 확인하며, 역할 조회가 불가능하거나 ADMIN이 아니면 메시지를 버리고 세션을 닫는다.
+ * 전달할 때 실제 물리 세션과 서버 내부의 불변 판정 결과를 대조한다. 판정 결과가 없거나 ADMIN이
+ * 아니면 메시지를 버린다. 판정 수집 뒤 연결된 세션은 닫지 않고, 후보였지만 ADMIN이 회수된
+ * 세션만 닫는다.
* RAG 개인 알림과 연결 제어 프레임은 이 검사 대상이 아니다.
*/
@Component
@RequiredArgsConstructor
-@Slf4j
public class StompDashboardOutboundAuthorizationInterceptor implements ChannelInterceptor {
private static final String DASHBOARD_TOPIC = "/topic/dashboard";
private final StompSessionRegistry stompSessionRegistry;
- private final RoleAuthorityService roleAuthorityService;
@Override
public Message> preSend(Message> message, MessageChannel channel) {
@@ -47,17 +44,19 @@ public Message> preSend(Message> message, MessageChannel channel) {
return null;
}
- // 3. 현재 primary의 역할을 확인한다. 조회 실패도 이전 ADMIN snapshot으로 우회하지 않는다.
- try {
- if (roleAuthorityService.getRolesForAdmin(authorization.userId()).contains("ADMIN")) {
- return message;
- }
- } catch (RuntimeException exception) {
- log.warn("STOMP 대시보드 전송의 최신 역할을 확인할 수 없어 차단합니다: {}", exception.getMessage());
+ // 3. 이 push의 primary 판정이 없으면 이전 CONNECT 역할로 폴백하지 않고 차단한다.
+ Object decision = message.getHeaders().get(DashboardAuthorizationSnapshot.HEADER);
+ if (!(decision instanceof DashboardAuthorizationSnapshot snapshot)) {
+ return null;
}
- // 4. 회수되었거나 확인할 수 없는 세션은 메시지를 버리고 다음 전송도 받지 않도록 닫는다.
- stompSessionRegistry.close(sessionId);
+ // 4. 늦게 연결된 세션은 이번 push만 건너뛰고, 확인된 비ADMIN 후보는 연결도 닫는다.
+ if (snapshot.allows(sessionId, authorization.userId())) {
+ return message;
+ }
+ if (snapshot.wasCandidate(sessionId, authorization.userId())) {
+ stompSessionRegistry.close(sessionId);
+ }
return null;
}
}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/dashboard/controller/DashboardWebSocketController.java b/backend/src/main/java/com/opensource/docgrid/domain/dashboard/controller/DashboardWebSocketController.java
index b3ee6c0c..033d6ea8 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/dashboard/controller/DashboardWebSocketController.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/dashboard/controller/DashboardWebSocketController.java
@@ -1,8 +1,18 @@
package com.opensource.docgrid.domain.dashboard.controller;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.stream.Collectors;
+
+import org.springframework.messaging.simp.SimpMessageHeaderAccessor;
+import org.springframework.messaging.simp.SimpMessageType;
import org.springframework.messaging.simp.SimpMessagingTemplate;
import org.springframework.stereotype.Component;
+import com.opensource.docgrid.domain.auth.service.query.PrimaryRoleQueryService;
+import com.opensource.docgrid.domain.auth.websocket.DashboardAuthorizationSnapshot;
+import com.opensource.docgrid.domain.auth.websocket.StompSessionRegistry;
import com.opensource.docgrid.domain.dashboard.dto.response.DashboardSummaryResponse;
import lombok.RequiredArgsConstructor;
@@ -12,7 +22,9 @@
*
*
이 컨트롤러는 직접 지표를 집계하지 않는다. 호출하는 쪽(재처리 트리거, 상태 전이 이벤트
* 리스너 등)이 {@code DashboardQueryService}로 최신 Snapshot을 계산해 넘겨주면 그대로
- * 브로드캐스트만 한다.
+ * 브로드캐스트한다. 호출자가 요약을 먼저 계산한 뒤 이 메서드에서 권한을 판정해야,
+ * 회수 전 판정으로 회수 후 새로 계산한 요약을 승인하지 않는다. 판정은 push 사이에 재사용하지 않는다.
+ * 구독자별 DB 조회를 반복하지 않도록 각 push의 ADMIN 후보를 primary에서 일괄 확인한다.
*/
@Component
@RequiredArgsConstructor
@@ -21,8 +33,30 @@ public class DashboardWebSocketController {
private static final String DASHBOARD_TOPIC = "/topic/dashboard";
private final SimpMessagingTemplate messagingTemplate;
+ private final StompSessionRegistry stompSessionRegistry;
+ private final PrimaryRoleQueryService primaryRoleQueryService;
public void sendDashboardUpdate(DashboardSummaryResponse summary) {
- messagingTemplate.convertAndSend(DASHBOARD_TOPIC, summary);
+ // 1. 이 백엔드의 열린 CONNECT-ADMIN 세션만 후보로 모으고 같은 사용자의 탭은 중복 제거한다.
+ Map candidateSessions = stompSessionRegistry.authenticatedSessions().stream()
+ .filter(session -> session.authorization().roles().contains("ADMIN"))
+ .collect(Collectors.toUnmodifiableMap(session -> session.sessionId(),
+ session -> session.authorization().userId()));
+ List candidateUserIds = candidateSessions.values().stream()
+ .distinct()
+ .toList();
+
+ // 2. 후보가 없으면 DB 트랜잭션도 열지 않는다. 실패 시에는 이전 권한으로 우회하지 않는다.
+ DashboardAuthorizationSnapshot snapshot = new DashboardAuthorizationSnapshot(candidateSessions,
+ candidateUserIds.isEmpty()
+ ? Set.of()
+ : primaryRoleQueryService.findCurrentAdminUserIds(candidateUserIds)
+ );
+
+ // 3. native STOMP 헤더가 아닌 서버 내부 헤더로만 이 push의 판정을 브로커에 전달한다.
+ SimpMessageHeaderAccessor headers = SimpMessageHeaderAccessor.create(SimpMessageType.MESSAGE);
+ headers.setHeader(DashboardAuthorizationSnapshot.HEADER, snapshot);
+ headers.setLeaveMutable(true);
+ messagingTemplate.convertAndSend(DASHBOARD_TOPIC, summary, headers.getMessageHeaders());
}
}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/dashboard/service/command/EmbeddingJobRetryService.java b/backend/src/main/java/com/opensource/docgrid/domain/dashboard/service/command/EmbeddingJobRetryService.java
index 5ad37df1..d44f884c 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/dashboard/service/command/EmbeddingJobRetryService.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/dashboard/service/command/EmbeddingJobRetryService.java
@@ -18,7 +18,8 @@
/**
* 관리자 재처리 버튼 클릭을 A 담당자의 {@link EmbeddingJobManualRetryService}로 위임하고,
- * 성공 시 최신 대시보드 집계를 WebSocket으로 push하는 Command Service.
+ * 성공 시 최신 대시보드 집계를 WebSocket으로 push하는 Command Service. 알림 인가 조회가
+ * 실패해도 이미 커밋된 재처리 결과를 HTTP 실패로 잘못 보고하지 않는다.
*
*
{@code embedding_jobs}를 직접 update하지 않는다 — 상태 전환은 전부 A의 Service를 경유한다.
* FAILED 목록 조회만 {@link EmbeddingJobRepository}를 직접 읽는다(쓰기가 아니므로 A/B 경계 위반이
@@ -49,7 +50,7 @@ public ManualRetriedIndexingJobResponse retryJob(Long jobId) {
// 1. 상태 전환은 A Service에 위임한다 — 여기서 예외가 나면(404/409) 그대로 전파시킨다.
ManualRetriedIndexingJobResponse response = embeddingJobManualRetryService.retry(jobId);
// 2. 재처리 성공 후에만 최신 집계를 다시 계산해서 push한다.
- dashboardWebSocketController.sendDashboardUpdate(dashboardQueryService.getSummary());
+ pushDashboardUpdateAfterRetry();
return response;
}
@@ -82,7 +83,7 @@ public RetryAllJobsResponse retryAllFailedJobs() {
// 3. 실제로 바뀐 게 있을 때만(1건 이상 성공) push한다 — 전부 실패하면 push할 변경사항이 없다.
if (retriedCount > 0) {
- dashboardWebSocketController.sendDashboardUpdate(dashboardQueryService.getSummary());
+ pushDashboardUpdateAfterRetry();
}
return new RetryAllJobsResponse(
@@ -93,4 +94,14 @@ public RetryAllJobsResponse retryAllFailedJobs() {
RETRY_ALL_MESSAGE_FORMAT.formatted(retriedCount, skippedCount, failedCount)
);
}
+
+ private void pushDashboardUpdateAfterRetry() {
+ try {
+ // 재처리는 이미 커밋됐다. 대시보드 알림만 실패하면 운영 결과를 실패로 뒤집지 않는다.
+ dashboardWebSocketController.sendDashboardUpdate(dashboardQueryService.getSummary());
+ } catch (RuntimeException exception) {
+ log.warn("재처리는 완료됐으나 대시보드 알림을 보내지 못했습니다. cause={}",
+ exception.getClass().getSimpleName());
+ }
+ }
}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/user/repository/UserRoleRepository.java b/backend/src/main/java/com/opensource/docgrid/domain/user/repository/UserRoleRepository.java
index ce131452..587b8b50 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/user/repository/UserRoleRepository.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/user/repository/UserRoleRepository.java
@@ -30,6 +30,10 @@ public interface UserRoleRepository extends JpaRepository {
@Query("SELECT ur.role.code FROM UserRole ur WHERE ur.user.id = :userId")
List findRoleCodesByUserId(@Param("userId") Long userId);
+ /** 한 대시보드 push의 후보 사용자 중 현재 ADMIN인 ID만 primary 트랜잭션에 반환한다. */
+ @Query("SELECT DISTINCT ur.user.id FROM UserRole ur WHERE ur.user.id IN :userIds AND ur.role.code = 'ADMIN'")
+ List findAdminUserIdsByUserIdIn(@Param("userIds") List userIds);
+
boolean existsByUserIdAndRoleCode(Long userId, String roleCode);
Optional findByUserIdAndRoleCode(Long userId, String roleCode);
diff --git a/backend/src/test/java/com/opensource/docgrid/domain/auth/integration/StompDashboardRealRoleRevocationIntegrationTest.java b/backend/src/test/java/com/opensource/docgrid/domain/auth/integration/StompDashboardRealRoleRevocationIntegrationTest.java
index dff7ff98..b8717519 100644
--- a/backend/src/test/java/com/opensource/docgrid/domain/auth/integration/StompDashboardRealRoleRevocationIntegrationTest.java
+++ b/backend/src/test/java/com/opensource/docgrid/domain/auth/integration/StompDashboardRealRoleRevocationIntegrationTest.java
@@ -2,14 +2,22 @@
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.ArgumentMatchers.anyList;
import java.lang.reflect.Type;
import java.time.Duration;
import java.time.LocalDateTime;
import java.util.UUID;
import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
@@ -20,12 +28,16 @@
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.boot.test.web.server.LocalServerPort;
import org.springframework.data.redis.core.StringRedisTemplate;
+import org.springframework.data.redis.core.script.RedisScript;
import org.springframework.http.HttpEntity;
import org.springframework.http.HttpHeaders;
import org.springframework.http.HttpMethod;
import org.springframework.http.HttpStatus;
import org.springframework.http.ResponseEntity;
+import org.springframework.messaging.Message;
import org.springframework.messaging.converter.MappingJackson2MessageConverter;
+import org.springframework.messaging.simp.SimpMessageHeaderAccessor;
+import org.springframework.messaging.simp.SimpMessageType;
import org.springframework.messaging.simp.stomp.StompCommand;
import org.springframework.messaging.simp.stomp.StompFrameHandler;
import org.springframework.messaging.simp.stomp.StompHeaders;
@@ -36,12 +48,15 @@
import org.springframework.test.context.ActiveProfiles;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
+import org.springframework.test.context.bean.override.mockito.MockitoSpyBean;
import org.springframework.boot.test.web.client.TestRestTemplate;
import org.springframework.web.socket.WebSocketHttpHeaders;
import org.springframework.web.socket.client.standard.StandardWebSocketClient;
import org.springframework.web.socket.messaging.WebSocketStompClient;
import com.opensource.docgrid.domain.auth.jwt.JwtProvider;
+import com.opensource.docgrid.domain.auth.websocket.DashboardAuthorizationSnapshot;
+import com.opensource.docgrid.domain.auth.websocket.StompDashboardOutboundAuthorizationInterceptor;
import com.opensource.docgrid.domain.dashboard.controller.DashboardWebSocketController;
import com.opensource.docgrid.domain.dashboard.dto.response.DashboardSummaryResponse;
import com.opensource.docgrid.domain.dashboard.dto.response.DocumentsSummaryResponse;
@@ -59,7 +74,9 @@
* 실제 HTTP 역할 회수, PostgreSQL 커밋, Redis 캐시 무효화와 STOMP 전송 인가를 연결한다.
*
*
테스트 전용 사용자만 생성·삭제하며, 자동 재검증을 늦춘 채 회수 직후의
- * 새 구독과 기존 구독 push가 차단되는지 확인한다. GCP 복제본 라우팅은 범위 밖이다.
+ * 새 구독과 기존 구독 push가 차단되는지 확인한다. 권한 확인을 마친 전송을 멈춰
+ * 회수 응답 뒤에 도착할 수 있는 기존 진행 중 메시지도 별도로 재현한다.
+ * GCP 복제본 라우팅은 범위 밖이다.
*/
@Tag("integration")
@ActiveProfiles("test")
@@ -80,7 +97,7 @@ class StompDashboardRealRoleRevocationIntegrationTest {
@Autowired
private TestRestTemplate restTemplate;
- @Autowired
+ @MockitoSpyBean
private StringRedisTemplate redisTemplate;
@Autowired
@@ -98,6 +115,9 @@ class StompDashboardRealRoleRevocationIntegrationTest {
@Autowired
private DashboardWebSocketController dashboardWebSocketController;
+ @MockitoSpyBean
+ private StompDashboardOutboundAuthorizationInterceptor outboundAuthorizationInterceptor;
+
private WebSocketStompClient stompClient;
private User adminCaller;
private User targetUser;
@@ -211,6 +231,165 @@ void revokesRealAdminRole_andBlocksDashboardSessionsBeforeScheduledRevalidation(
assertThat(redisTemplate.opsForValue().get(cacheKey)).isEqualTo("ADMIN");
}
+ @Test
+ @DisplayName("권한 확인을 마친 진행 중 메시지는 회수 응답 뒤에도 도착할 수 있다")
+ void inFlightMessage_canArriveAfterRevocationResponse_whenAuthorizationFinishedFirst() throws Exception {
+ // 1. 실제 DB 역할과 WebSocket 구독을 준비하고, 전송 직전의 대기 지점만 테스트용 spy로 제어한다.
+ Role adminRole = roleRepository.findByCode("ADMIN").orElseThrow();
+ adminCaller = createTestUser("race-caller");
+ targetUser = createTestUser("race-target");
+ grantAdmin(adminCaller, adminRole);
+ grantAdmin(targetUser, adminRole);
+ BlockingQueue received = new LinkedBlockingQueue<>();
+ oldSession = connect(
+ jwtProvider.generateToken(targetUser.getId(), targetUser.getEmail()),
+ new LinkedBlockingQueue<>()
+ );
+ oldSession.subscribe(DASHBOARD_TOPIC, dashboardFrames(received));
+ await().atMost(Duration.ofSeconds(TIMEOUT_SECONDS))
+ .until(() -> hasDashboardSubscription(targetUser.getEmail()));
+
+ CountDownLatch authorizationFinished = new CountDownLatch(1);
+ CountDownLatch allowDelivery = new CountDownLatch(1);
+ AtomicBoolean holdOneMessage = new AtomicBoolean(true);
+ doAnswer(invocation -> {
+ Message> authorized = (Message>) invocation.callRealMethod();
+ if (authorized != null
+ && SimpMessageType.MESSAGE.equals(SimpMessageHeaderAccessor.getMessageType(authorized.getHeaders()))
+ && DASHBOARD_TOPIC.equals(SimpMessageHeaderAccessor.getDestination(authorized.getHeaders()))
+ && holdOneMessage.compareAndSet(true, false)) {
+ authorizationFinished.countDown();
+ if (!allowDelivery.await(TIMEOUT_SECONDS, TimeUnit.SECONDS)) {
+ throw new IllegalStateException("테스트 전송 대기 시간이 초과됐습니다.");
+ }
+ }
+ return authorized;
+ }).when(outboundAuthorizationInterceptor).preSend(any(), any());
+
+ CompletableFuture dispatch = CompletableFuture.runAsync(
+ () -> dashboardWebSocketController.sendDashboardUpdate(sampleSummary(42L))
+ );
+ try {
+ assertThat(authorizationFinished.await(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isTrue();
+ assertThat(received.poll(NO_DELIVERY_MILLIS, TimeUnit.MILLISECONDS)).isNull();
+
+ // 2. 이미 ADMIN으로 허용된 메시지를 붙잡은 동안 실제 HTTP 회수·DB 커밋을 완료한다.
+ HttpHeaders headers = new HttpHeaders();
+ headers.setBearerAuth(jwtProvider.generateToken(adminCaller.getId(), adminCaller.getEmail()));
+ ResponseEntity response = restTemplate.exchange(
+ "/admin/users/" + targetUser.getId() + "/roles/ADMIN",
+ HttpMethod.DELETE,
+ new HttpEntity<>(null, headers),
+ String.class
+ );
+ assertThat(response.getStatusCode()).isEqualTo(HttpStatus.OK);
+ assertThat(userRoleRepository.existsByUserIdAndRoleCode(targetUser.getId(), "ADMIN")).isFalse();
+ assertThat(received.poll(NO_DELIVERY_MILLIS, TimeUnit.MILLISECONDS)).isNull();
+
+ // 3. 회수 응답 뒤에 전송을 풀어, 새로 허용된 전송과 기존 진행 중 전송을 구별한다.
+ allowDelivery.countDown();
+ dispatch.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ DashboardSummaryResponse inFlightMessage = received.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ assertThat(inFlightMessage).isNotNull();
+ assertThat(inFlightMessage.documents().total()).isEqualTo(42L);
+ } finally {
+ allowDelivery.countDown();
+ dispatch.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ }
+ }
+
+ @Test
+ @DisplayName("Redis 무효화 실패로 낡은 ADMIN 캐시가 남아도 새 대시보드 push를 차단한다")
+ @SuppressWarnings("unchecked")
+ void redisInvalidationFailure_doesNotAuthorizeRevokedDashboardSession() throws Exception {
+ // 1. 실제 구독과 캐시를 만들고 역할 무효화 Lua 호출만 시험에서 실패시킨다.
+ Role adminRole = roleRepository.findByCode("ADMIN").orElseThrow();
+ adminCaller = createTestUser("redis-failure-caller");
+ targetUser = createTestUser("redis-failure-target");
+ grantAdmin(adminCaller, adminRole);
+ grantAdmin(targetUser, adminRole);
+ BlockingQueue received = new LinkedBlockingQueue<>();
+ oldSession = connect(
+ jwtProvider.generateToken(targetUser.getId(), targetUser.getEmail()),
+ new LinkedBlockingQueue<>()
+ );
+ oldSession.subscribe(DASHBOARD_TOPIC, dashboardFrames(received));
+ await().atMost(Duration.ofSeconds(TIMEOUT_SECONDS))
+ .until(() -> hasDashboardSubscription(targetUser.getEmail()));
+ redisTemplate.opsForValue().set(roleCacheKey(targetUser.getId()), "ADMIN", Duration.ofSeconds(30));
+ doThrow(new IllegalStateException("simulated Redis invalidate failure"))
+ .when(redisTemplate).execute(any(RedisScript.class), anyList());
+
+ // 2. DB 커밋과 HTTP 200은 완료되지만 Redis의 이전 ADMIN 캐시는 그대로 남는다.
+ HttpHeaders headers = new HttpHeaders();
+ headers.setBearerAuth(jwtProvider.generateToken(adminCaller.getId(), adminCaller.getEmail()));
+ ResponseEntity response = restTemplate.exchange(
+ "/admin/users/" + targetUser.getId() + "/roles/ADMIN",
+ HttpMethod.DELETE,
+ new HttpEntity<>(null, headers),
+ String.class
+ );
+ assertThat(response.getStatusCode()).isEqualTo(HttpStatus.OK);
+ assertThat(userRoleRepository.existsByUserIdAndRoleCode(targetUser.getId(), "ADMIN")).isFalse();
+ assertThat(redisTemplate.opsForValue().get(roleCacheKey(targetUser.getId()))).isEqualTo("ADMIN");
+
+ // 3. 다음 push는 Redis가 아닌 primary 일괄 조회를 사용하므로 메시지를 버린다.
+ dashboardWebSocketController.sendDashboardUpdate(sampleSummary(43L));
+ assertThat(received.poll(NO_DELIVERY_MILLIS, TimeUnit.MILLISECONDS)).isNull();
+ await().atMost(Duration.ofSeconds(TIMEOUT_SECONDS)).until(() -> !oldSession.isConnected());
+ }
+
+ @Test
+ @DisplayName("실제 발행 코드의 primary 판정은 outbound에 남고 STOMP 클라이언트에는 노출되지 않는다")
+ void serverOnlyHeader_reachesOutboundButNotClient() throws Exception {
+ // 1. 실제 구독을 만들고 수신자의 프레임 헤더와 서버 outbound 헤더를 따로 관측한다.
+ Role adminRole = roleRepository.findByCode("ADMIN").orElseThrow();
+ targetUser = createTestUser("internal-header");
+ grantAdmin(targetUser, adminRole);
+ BlockingQueue clientHeaders = new LinkedBlockingQueue<>();
+ BlockingQueue received = new LinkedBlockingQueue<>();
+ oldSession = connect(
+ jwtProvider.generateToken(targetUser.getId(), targetUser.getEmail()),
+ new LinkedBlockingQueue<>()
+ );
+ oldSession.subscribe(DASHBOARD_TOPIC, new StompFrameHandler() {
+ @Override
+ public Type getPayloadType(StompHeaders headers) {
+ return DashboardSummaryResponse.class;
+ }
+
+ @Override
+ public void handleFrame(StompHeaders headers, Object payload) {
+ clientHeaders.add(headers);
+ received.add((DashboardSummaryResponse) payload);
+ }
+ });
+ await().atMost(Duration.ofSeconds(TIMEOUT_SECONDS))
+ .until(() -> hasDashboardSubscription(targetUser.getEmail()));
+
+ AtomicReference