Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
package com.opensource.docgrid.domain.auth.websocket;

import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.simp.SimpMessageHeaderAccessor;
import org.springframework.messaging.simp.SimpMessageType;
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 역할로 제한한다.
*
* <p>SUBSCRIBE 이후 역할이 회수되어도 브로커에는 구독이 남을 수 있다. 이 경계는 각 물리 세션에
* 전달할 때 다시 확인하며, 역할 조회가 불가능하거나 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) {
// 1. 브로커가 구독자에게 보내는 대시보드 MESSAGE만 검사한다.
if (!SimpMessageType.MESSAGE.equals(SimpMessageHeaderAccessor.getMessageType(message.getHeaders()))
|| !DASHBOARD_TOPIC.equals(SimpMessageHeaderAccessor.getDestination(message.getHeaders()))) {
return message;
}

// 2. 실제 수신 세션의 인증 snapshot이 없거나 CONNECT 당시 ADMIN이 아니면 버린다.
String sessionId = SimpMessageHeaderAccessor.getSessionId(message.getHeaders());
StompSessionAuthorization authorization = stompSessionRegistry.authorizationFor(sessionId);
if (authorization == null || !authorization.roles().contains("ADMIN")) {
if (sessionId != null) {
stompSessionRegistry.close(sessionId);
}
return null;
}

// 3. 현재 primary의 역할을 확인한다. 조회 실패도 이전 ADMIN snapshot으로 우회하지 않는다.
try {
if (roleAuthorityService.getRolesForAdmin(authorization.userId()).contains("ADMIN")) {
return message;
}
} catch (RuntimeException exception) {
log.warn("STOMP 대시보드 전송의 최신 역할을 확인할 수 없어 차단합니다: {}", exception.getMessage());
}

// 4. 회수되었거나 확인할 수 없는 세션은 메시지를 버리고 다음 전송도 받지 않도록 닫는다.
stompSessionRegistry.close(sessionId);
return null;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,12 +14,18 @@
import org.springframework.security.core.GrantedAuthority;
import org.springframework.stereotype.Component;

import com.opensource.docgrid.domain.auth.jwt.RoleAuthorityService;

import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;

/**
* STOMP client의 SUBSCRIBE·SEND 목적지를 명시적인 허용 목록으로 제한한다.
*
* <p>{@code StompAuthChannelInterceptor}가 CONNECT 시점에 세션에 부착한 Principal을 재사용해
* 목적지 접근 시점에 다시 검증한다. {@code /topic/dashboard}는 ADMIN만, 사용자별 RAG 완료 알림인
* {@code /user/queue/rag-answer}는 인증된 사용자만 구독할 수 있다. 그 외 정확한 목적지와 pattern,
* <p>{@code StompAuthChannelInterceptor}가 CONNECT 시점에 세션에 부착한 Principal을 확인하고,
* 관리자 구독은 현재 primary의 역할도 다시 조회한다. {@code /topic/dashboard}는 ADMIN만,
* 사용자별 RAG 완료 알림인 {@code /user/queue/rag-answer}는 인증된 사용자만 구독할 수 있다.
* 그 외 정확한 목적지와 pattern,
* Spring이 내부에서 만드는 실제 {@code /queue} 목적지는 모두 거부한다.
*
* <p>애플리케이션에는 client가 호출할 {@code @MessageMapping}이 없고 실제 push는 서버의
Expand All @@ -32,6 +38,8 @@
* {@code MissingCsrfTokenException}으로 거부되므로, 그 DSL 대신 이 수동 Interceptor로 구현한다.
*/
@Component
@RequiredArgsConstructor
@Slf4j
public class StompDestinationAuthorizationInterceptor implements ChannelInterceptor {

private static final String DASHBOARD_TOPIC = "/topic/dashboard";
Expand All @@ -40,6 +48,8 @@ public class StompDestinationAuthorizationInterceptor implements ChannelIntercep
private static final String SUBSCRIPTION_DENIED_MESSAGE = "구독 권한이 없습니다.";
private static final String SEND_DENIED_MESSAGE = "메시지를 보낼 권한이 없습니다.";

private final RoleAuthorityService roleAuthorityService;

@Override
public Message<?> preSend(Message<?> message, MessageChannel channel) {
StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class);
Expand All @@ -59,7 +69,7 @@ public Message<?> preSend(Message<?> message, MessageChannel channel) {

// 3. 넓은 pattern 대신 프런트가 실제 사용하는 두 목적지만 정확히 일치할 때 허용한다.
String destination = accessor.getDestination();
if (DASHBOARD_TOPIC.equals(destination) && isAdmin(accessor.getUser())) {
if (DASHBOARD_TOPIC.equals(destination) && isCurrentAdmin(accessor.getUser())) {
return message;
}
if (RAG_ANSWER_QUEUE.equals(destination) && isAuthenticated(accessor.getUser())) {
Expand All @@ -70,13 +80,24 @@ public Message<?> preSend(Message<?> message, MessageChannel channel) {
throw new AccessDeniedException(SUBSCRIPTION_DENIED_MESSAGE);
}

private boolean isAdmin(Principal user) {
private boolean isCurrentAdmin(Principal user) {
if (!(user instanceof Authentication authentication)) {
return false;
}
return authentication.isAuthenticated() && authentication.getAuthorities().stream()
// CONNECT 때의 ADMIN snapshot만으로는 역할 회수 뒤 새 구독을 허용할 수 없다.
boolean wasAdmin = authentication.isAuthenticated() && authentication.getAuthorities().stream()
.map(GrantedAuthority::getAuthority)
.anyMatch(ADMIN_AUTHORITY::equals);
if (!wasAdmin || !(authentication.getDetails() instanceof Long userId)) {
return false;
}
try {
return roleAuthorityService.getRolesForAdmin(userId).contains("ADMIN");
} catch (RuntimeException exception) {
// primary 상태를 확인하지 못하면 저장된 권한으로 폴백하지 않는다.
log.warn("STOMP 관리자 구독의 최신 역할을 확인할 수 없어 거부합니다: {}", exception.getMessage());
return false;
}
}

private boolean isAuthenticated(Principal user) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,8 @@
* 현재 Backend 인스턴스가 소유한 물리 WebSocket 연결과 STOMP 인증 snapshot을 함께 관리한다.
*
* <p>물리 연결은 WebSocket decorator가 먼저 등록하고, CONNECT 인증이 끝난 뒤 같은 sessionId에
* 인증 snapshot을 결합한다. 주기 검사는 인증 완료 세션만 읽으며, 종료와 인증이 경합해도 하나의
* 인증 snapshot을 결합한다. 주기 검사와 outbound 인가는 인증 완료 세션만 읽으며,
* 종료와 인증이 경합해도 하나의
* ConcurrentMap entry를 기준으로 정리해 닫힌 연결이 다시 등록되는 것을 막는다.
*
* <p>이 registry는 로컬 전송 자원만 관리한다. 여러 Backend 인스턴스는 각자 자신의 registry를
Expand Down Expand Up @@ -69,6 +70,15 @@ public int authenticatedSessionCount() {
.count();
}

/** outbound 전송 대상의 열린 물리 세션에 결합된 CONNECT 인증 snapshot을 조회한다. */
public StompSessionAuthorization authorizationFor(String sessionId) {
if (sessionId == null) {
return null;
}
SessionState state = sessions.get(sessionId);
return state != null && state.session().isOpen() ? state.authorization() : null;
}

public boolean close(String sessionId) {
SessionState state = sessions.get(sessionId);
if (state == null) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@

import com.opensource.docgrid.domain.auth.jwt.StompAuthChannelInterceptor;
import com.opensource.docgrid.domain.auth.websocket.StompDestinationAuthorizationInterceptor;
import com.opensource.docgrid.domain.auth.websocket.StompDashboardOutboundAuthorizationInterceptor;
import com.opensource.docgrid.domain.auth.websocket.StompSessionTrackingDecoratorFactory;

import lombok.RequiredArgsConstructor;
Expand All @@ -19,8 +20,9 @@
*
* <p>인증·인가는 이 설정이 아니라 {@link StompAuthChannelInterceptor}(CONNECT 시점 인증)와
* {@link StompDestinationAuthorizationInterceptor}(SUBSCRIBE·SEND 시점 인가)가 담당한다.
* 이 클래스는 전송 계층 구성(endpoint·broker·origin), 물리 세션 추적과 두 Interceptor의 등록 순서만
* 책임진다.
* 기존 구독에 나가는 대시보드 메시지는 {@link StompDashboardOutboundAuthorizationInterceptor}가
* 현재 역할을 다시 확인한다. 이 클래스는 전송 계층 구성(endpoint·broker·origin), 물리 세션 추적과
* Interceptor 등록만 책임진다.
*
* <p>{@code /queue}는 RAG 답변 개인 알림({@code convertAndSendToUser})의 broker 내부 목적지로 쓰인다.
* client는 {@code /user/queue/rag-answer}만 구독할 수 있고, 실제 {@code /queue}와 pattern 접근은
Expand All @@ -33,6 +35,7 @@ public class WebSocketConfig implements WebSocketMessageBrokerConfigurer {

private final StompAuthChannelInterceptor stompAuthChannelInterceptor;
private final StompDestinationAuthorizationInterceptor stompDestinationAuthorizationInterceptor;
private final StompDashboardOutboundAuthorizationInterceptor stompDashboardOutboundAuthorizationInterceptor;
private final StompSessionTrackingDecoratorFactory stompSessionTrackingDecoratorFactory;

@Override
Expand Down Expand Up @@ -60,4 +63,10 @@ public void configureClientInboundChannel(ChannelRegistration registration) {
// 순서가 바뀌면 2번 시점에 Principal이 아직 없어 항상 거부된다.
registration.interceptors(stompAuthChannelInterceptor, stompDestinationAuthorizationInterceptor);
}

@Override
public void configureClientOutboundChannel(ChannelRegistration registration) {
// 브로커가 기존 구독자에게 보내는 매 메시지는 실제 전달 전에 최신 관리자 역할을 확인한다.
registration.interceptors(stompDashboardOutboundAuthorizationInterceptor);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,6 @@
import org.springframework.web.socket.messaging.WebSocketStompClient;

import com.opensource.docgrid.domain.auth.jwt.JwtProvider;
import com.opensource.docgrid.domain.auth.websocket.StompSessionRevalidationScheduler;
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;
Expand All @@ -57,10 +56,10 @@
import com.opensource.docgrid.domain.user.repository.UserRoleRepository;

/**
* 실제 HTTP 역할 회수, PostgreSQL 커밋, Redis 캐시 무효화와 STOMP 세션 재검증을 연결한다.
* 실제 HTTP 역할 회수, PostgreSQL 커밋, Redis 캐시 무효화와 STOMP 전송 인가를 연결한다.
*
* <p>테스트 전용 사용자만 생성·삭제하며, 자동 재검증을 늦추고 직접 호출해 회수 직후의
* 기존 구독 창과 재검증 후 종료를 구분한다. GCP 복제본 라우팅은 이 로컬 테스트의 범위가 아니다.
* <p>테스트 전용 사용자만 생성·삭제하며, 자동 재검증을 늦춘 채 회수 직후의
* 새 구독과 기존 구독 push가 차단되는지 확인한다. GCP 복제본 라우팅은 범위 밖이다.
*/
@Tag("integration")
@ActiveProfiles("test")
Expand Down Expand Up @@ -96,16 +95,15 @@ class StompDashboardRealRoleRevocationIntegrationTest {
@Autowired
private SimpUserRegistry simpUserRegistry;

@Autowired
private StompSessionRevalidationScheduler revalidationScheduler;

@Autowired
private DashboardWebSocketController dashboardWebSocketController;

private WebSocketStompClient stompClient;
private User adminCaller;
private User targetUser;
private StompSession adminSession;
private StompSession oldSession;
private StompSession oldIdleSession;
private StompSession newSession;

@DynamicPropertySource
Expand All @@ -126,21 +124,25 @@ void setUp() {
void tearDown() {
// 1. 세션을 먼저 닫고 이 테스트가 만든 사용자·역할과 Redis 키만 제거한다.
disconnect(newSession);
disconnect(oldIdleSession);
disconnect(oldSession);
disconnect(adminSession);
if (targetUser != null) {
redisTemplate.delete(roleCacheKey(targetUser.getId()));
redisTemplate.delete(roleEpochKey(targetUser.getId()));
deleteTestUser(targetUser.getId());
}
if (adminCaller != null) {
redisTemplate.delete(roleCacheKey(adminCaller.getId()));
redisTemplate.delete(roleEpochKey(adminCaller.getId()));
deleteTestUser(adminCaller.getId());
}
stompClient.stop();
}

@Test
@DisplayName("HTTP 역할 회수 커밋이 Redis 캐시와 새·기존 관리자 WebSocket 세션에 반영된다")
void revokesRealAdminRole_andRevalidatesDashboardSessions() throws Exception {
@DisplayName("HTTP 역할 회수 뒤 주기 재검증 없이 새 구독과 기존 구독 push가 차단된다")
void revokesRealAdminRole_andBlocksDashboardSessionsBeforeScheduledRevalidation() throws Exception {
// 1. seed 역할은 재사용하되 사용자와 매핑은 이 테스트만 소유한다.
Role adminRole = roleRepository.findByCode("ADMIN").orElseThrow();
adminCaller = createTestUser("caller");
Expand All @@ -156,10 +158,21 @@ void revokesRealAdminRole_andRevalidatesDashboardSessions() throws Exception {

// 2. 실제 STOMP CONNECT와 SUBSCRIBE가 DB 역할을 Redis에 캐시하고 브로커에 등록한다.
BlockingQueue<DashboardSummaryResponse> received = new LinkedBlockingQueue<>();
BlockingQueue<DashboardSummaryResponse> activeReceived = new LinkedBlockingQueue<>();
oldSession = connect(targetToken, new LinkedBlockingQueue<>());
oldSession.subscribe(DASHBOARD_TOPIC, dashboardFrames(received));
adminSession = connect(
jwtProvider.generateToken(adminCaller.getId(), adminCaller.getEmail()),
new LinkedBlockingQueue<>()
);
adminSession.subscribe(DASHBOARD_TOPIC, dashboardFrames(activeReceived));
await().atMost(Duration.ofSeconds(TIMEOUT_SECONDS))
.until(() -> hasDashboardSubscription(targetEmail));
.until(() -> hasDashboardSubscription(targetEmail) && hasDashboardSubscription(adminCaller.getEmail()));
dashboardWebSocketController.sendDashboardUpdate(sampleSummary(1L));
assertThat(received.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isNotNull();
assertThat(activeReceived.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isNotNull();
BlockingQueue<Throwable> oldIdleFailures = new LinkedBlockingQueue<>();
oldIdleSession = connect(targetToken, oldIdleFailures);
assertThat(redisTemplate.opsForValue().get(cacheKey)).contains("ADMIN");

// 3. 실제 관리자 HTTP 요청이 DB 트랜잭션을 커밋한 뒤 Redis 캐시를 무효화한다.
Expand All @@ -176,24 +189,26 @@ void revokesRealAdminRole_andRevalidatesDashboardSessions() throws Exception {
assertThat(redisTemplate.opsForValue().get(cacheKey)).isNull();
assertThat(redisTemplate.opsForValue().get(epochKey)).isEqualTo("1");

// 4. 주기 검사 전 옛 구독에는 push가 도달하지만 새 연결의 관리자 구독은 거부된다.
dashboardWebSocketController.sendDashboardUpdate(sampleSummary(1L));
assertThat(received.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isNotNull();
// 실제 무효화 결과를 확인한 뒤 의도적으로 옛 ADMIN 캐시를 넣어도 관리자 경계는 primary만 신뢰한다.
redisTemplate.opsForValue().set(cacheKey, "ADMIN", Duration.ofSeconds(30));

// 4. 회수 뒤 기존 연결의 새 SUBSCRIBE와 기존 구독의 새 push가 모두 차단된다.
oldIdleSession.subscribe(DASHBOARD_TOPIC, dashboardFrames(new LinkedBlockingQueue<>()));
assertThat(oldIdleFailures.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isNotNull();
dashboardWebSocketController.sendDashboardUpdate(sampleSummary(2L));
await().atMost(Duration.ofSeconds(TIMEOUT_SECONDS))
.until(() -> !oldSession.isConnected());
assertThat(received.poll(NO_DELIVERY_MILLIS, TimeUnit.MILLISECONDS)).isNull();
assertThat(activeReceived.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isNotNull();

// 5. 새 CONNECT도 회수된 관리자 구독을 만들 수 없고 scheduler를 호출할 필요가 없다.
BlockingQueue<Throwable> newFailures = new LinkedBlockingQueue<>();
newSession = connect(targetToken, newFailures);
newSession.subscribe(DASHBOARD_TOPIC, dashboardFrames(new LinkedBlockingQueue<>()));
assertThat(newFailures.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isNotNull();
await().atMost(Duration.ofSeconds(TIMEOUT_SECONDS))
.until(() -> !newSession.isConnected());
assertThat(redisTemplate.opsForValue().get(cacheKey)).doesNotContain("ADMIN");

// 5. 실제 DB를 읽는 재검증 뒤 옛 세션과 구독이 제거되고 새 push는 전달되지 않는다.
revalidationScheduler.revalidate();
await().atMost(Duration.ofSeconds(TIMEOUT_SECONDS))
.until(() -> simpUserRegistry.getUser(targetEmail) == null);
assertThat(oldSession.isConnected()).isFalse();
dashboardWebSocketController.sendDashboardUpdate(sampleSummary(2L));
assertThat(received.poll(NO_DELIVERY_MILLIS, TimeUnit.MILLISECONDS)).isNull();
assertThat(redisTemplate.opsForValue().get(cacheKey)).isEqualTo("ADMIN");
}

private User createTestUser(String kind) {
Expand Down
Loading
Loading