From 3f1a3fcd6fd1603bfee04b3d24faa8fdf41e6418 Mon Sep 17 00:00:00 2001 From: kangcheolung Date: Mon, 10 Aug 2026 18:52:20 +0900 Subject: [PATCH 1/9] =?UTF-8?q?docs:=20#137=20RAGOps=20Dashboard=20WebSock?= =?UTF-8?q?et=20=EC=8B=A4=EC=8B=9C=EA=B0=84=20push=20=EC=84=A4=EA=B3=84=20?= =?UTF-8?q?=EB=AC=B8=EC=84=9C=20=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Sonnet 5 --- ...ng-#137-ragops-dashboard-websocket-push.md | 395 ++++++++++++++++++ 1 file changed, 395 insertions(+) create mode 100644 docs/design/kangcheolung-#137-ragops-dashboard-websocket-push.md diff --git a/docs/design/kangcheolung-#137-ragops-dashboard-websocket-push.md b/docs/design/kangcheolung-#137-ragops-dashboard-websocket-push.md new file mode 100644 index 0000000..0bce997 --- /dev/null +++ b/docs/design/kangcheolung-#137-ragops-dashboard-websocket-push.md @@ -0,0 +1,395 @@ +# Issue #137 RAGOps Dashboard WebSocket 실시간 push 상세 설계 + +closes #137 + +## 1. 배경과 목적 + +이슈#134(RAGOps Dashboard 집계 지표 조회 API)는 관리자가 요청할 때마다 스냅샷을 계산해 반환하는 +정적 조회만 제공한다. 이번 작업은 그 스냅샷을 상태 변경 시점에 서버가 먼저 push하는 실시간 +채널을 추가한다. Polling·고정 주기 스케줄러 방식은 채택하지 않고, 상태 변경 이벤트가 발생하는 +시점에만 push하는 이벤트 기반 구조로 간다. + +이번 이슈는 push 인프라(STOMP endpoint, 인증·인가, 전송 계층)까지만 만든다. 실제로 언제 +push를 트리거할지(재처리 클릭, Worker 상태 전이)는 후속 이슈가 이 인프라를 재사용해서 채운다. + +### 1.1 성공 기준 + +- `/ws`(SockJS) 연결과 `/topic/dashboard` 구독이 정상 동작한다. +- 상태 변경 시점에만 push하며, 고정 주기 폴링을 쓰지 않는다. +- ADMIN이 아닌 사용자는 연결 또는 구독 시점에 거부된다. +- 이 앱의 JWT 인증(stateless, 세션 없음) 모델과 충돌 없이 동작한다. + +## 2. 범위 + +### 2.1 포함 + +- STOMP endpoint(`/ws`)와 Message Broker(`/topic`) 설정 +- WebSocket 세션에 대한 JWT 인증(CONNECT 시점) 및 목적지별 인가(SUBSCRIBE 시점) +- `DashboardWebSocketController.sendDashboardUpdate()` — 수동 트리거 기반 push 전송 계층 +- 단위·통합 테스트 + +### 2.2 제외 + +- 실제 push 트리거 시점(재처리 API, Worker 상태 전이 이벤트 훅) — 후속 이슈 +- FAILED 작업 목록 조회, 재처리 API — 후속 이슈 +- Dashboard 화면(Frontend) + +## 3. 프로토콜 계약 + +``` +WebSocket 연결: /ws (SockJS) +구독 채널: /topic/dashboard +``` + +- HTTP 핸드셰이크(`/ws/**`)는 인증하지 않는다(permitAll). 인증·인가는 STOMP 레이어에서 처리한다. +- STOMP endpoint의 허용 origin은 `CorsConfig.ALLOWED_ORIGINS`를 재사용한다(REST API와 동일 목록). + +## 4. 인증·인가 설계 + +### 4.1 왜 HTTP 핸드셰이크가 아니라 STOMP 프레임에서 인증하는가 + +이 앱의 JWT 인증은 `Authorization` HTTP 헤더 기반이다. 그런데 SockJS가 `websocket` Transport로 +직접 연결하면 브라우저 네이티브 `WebSocket()`은 Upgrade 요청에 커스텀 헤더를 실을 수 없다. +`xhr-streaming`/`xhr-polling` 폴백에서는 되고 `websocket` Transport에서는 안 되는, Transport +종류에 따라 인증 여부가 갈리는 상태가 된다. + +그래서 인증을 HTTP 핸드셰이크가 아니라 **STOMP CONNECT 프레임**으로 옮겼다. STOMP 프로토콜은 +Transport와 무관하게 항상 커스텀 헤더(native header)를 실을 수 있다. `SecurityConfig`는 +`/ws/**`를 `permitAll()`로 열어 HTTP 레벨 인증을 하지 않고, 대신: + +- **CONNECT 시점 인증** — `StompAuthChannelInterceptor`가 CONNECT 프레임의 `Authorization` + 헤더로 JWT를 검증하고, 유효하면 WebSocket 세션에 `Authentication`을 부착한다(`accessor.setUser()`). + 토큰이 없거나 무효하면 그 자리에서 연결을 끊는다. +- **SUBSCRIBE 시점 인가** — `DashboardSubscriptionAuthorizationInterceptor`가 `/topic/dashboard` + 구독 요청마다 세션에 부착된 Principal이 `ROLE_ADMIN`인지 다시 확인한다. CONNECT 검증 하나에만 + 의존하지 않는 이중 방어다. + +### 4.2 `@EnableWebSocketSecurity`를 채택하지 않은 이유 + +원래는 SUBSCRIBE 인가를 Spring Security의 `@EnableWebSocketSecurity`(`MessageMatcherDelegatingAuthorizationManager` +DSL)로 구현하려 했다. 그런데 실제로 붙여보니 유효한 ADMIN 토큰으로도 모든 CONNECT가 +`MissingCsrfTokenException`으로 거부됐다. + +원인을 바이트코드 레벨까지 확인한 결과, `@EnableWebSocketSecurity`는 STOMP endpoint(`stompWebSocketHandlerMapping` +빈)가 등록된 걸 감지하면 `CsrfChannelInterceptor`를 무조건 함께 등록한다. 이 인터셉터는 +CONNECT 프레임 처리 시 세션에 저장된 `CsrfToken`이 있는지를 확인하는데, 이 앱은 +`SessionCreationPolicy.STATELESS` + HTTP CSRF 비활성화 상태라 세션 자체가 없어 `CsrfToken`이 +존재할 수 없다. 즉 ADMIN 여부와 무관하게 모든 CONNECT가 거부되는 구조였다. 이 버전(Spring +Security 6.5.11)에는 이 CSRF 요구를 끄는 공개 API가 없어서, `@EnableWebSocketSecurity` 자체를 +쓰지 않고 `DashboardSubscriptionAuthorizationInterceptor`를 순수 `ChannelInterceptor`로 직접 +구현했다. 그 결과 `spring-security-messaging` 의존성도 필요 없어져서 제외했다. + +### 4.3 구현 시 발견한 버그: `StompHeaderAccessor.wrap()` vs `MessageHeaderAccessor.getAccessor()` + +`StompAuthChannelInterceptor`에서 처음엔 `StompHeaderAccessor.wrap(message)`로 Accessor를 +가져와 `setUser()`를 호출했는데, `wrap()`은 검증 전용 복사본을 만들 뿐이라 그 위에서 호출한 +`setUser()` 결과가 반환되는 `message`에 반영되지 않았다. 그 결과 SUBSCRIBE 단계에서 +`accessor.getUser()`가 항상 `null`이었다. `StompSubProtocolHandler`가 `leaveMutable(true)`로 +캐시해 둔 원본 Accessor를 가져오는 `MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class)`로 +바꿔서 해결했다. 실제 STOMP Client로 통합 테스트를 돌리기 전까지는 드러나지 않던 문제였다. + +## 5. 구현 상세 + +읽는 순서는 설정(5.1~5.3) → 인증(5.4) → 인가(5.5) → HTTP 레벨 배선(5.6) → 실제 push 전송(5.7) +순서를 따른다. 각 파일이 이전 파일의 결과를 어떻게 이어받는지가 핵심이다. + +### 5.1 `build.gradle` — WebSocket 의존성 추가 + +```diff + implementation 'org.springframework.boot:spring-boot-starter-security' + implementation 'org.springframework.boot:spring-boot-starter-validation' + implementation 'org.springframework.boot:spring-boot-starter-web' ++ implementation 'org.springframework.boot:spring-boot-starter-websocket' + implementation 'org.springframework.ai:spring-ai-starter-mcp-server-webmvc' +``` + +한 줄 추가. 이게 없으면 `@EnableWebSocketMessageBroker`, `StompHeaderAccessor`, `SimpMessagingTemplate` +등 이 이슈에서 쓰는 클래스가 클래스패스에 없어서 5.3부터 전부 컴파일이 안 된다. 프로젝트에 +WebSocket 관련 코드가 지금까지 전혀 없었기 때문에 순수 신규 추가다. + +`spring-security-messaging`은 한때 추가했다가 뺐다 — SUBSCRIBE 인가를 Spring Security의 +`@EnableWebSocketSecurity` DSL로 구현하려고 추가했었는데, 4.2절에서 설명하는 CSRF 문제 때문에 +그 방식을 버리면서 이 의존성도 같이 뺐다. 지금 `build.gradle`엔 흔적이 없다. + +### 5.2 `CorsConfig.java` — origin 목록을 WebSocket 설정과 공유 + +```diff + @Configuration + public class CorsConfig implements WebMvcConfigurer { + +- private static final List ALLOWED_ORIGINS = List.of( ++ // WebSocketConfig가 STOMP endpoint 허용 origin으로 재사용하므로 package-private으로 둔다. ++ static final List ALLOWED_ORIGINS = List.of( + "http://localhost:3000", + "http://localhost:8080" + ); +``` + +값은 안 바꾸고 접근 제어자만 `private` → package-private(제어자 없음)으로 넓혔다. REST API의 +CORS 설정과 STOMP endpoint의 origin 허용 설정은 스프링에서 서로 완전히 다른 설정 지점이라, +`WebSocketConfig`도 `/ws` 연결을 허용할 origin이 따로 필요하다. 이걸 또 하드코딩하면 값이 +두 파일에 중복돼서 배포 도메인이 바뀔 때 한쪽만 고치고 잊어버리는 사고가 날 수 있다. 같은 +패키지(`global.config`)에서만 쓰면 되니 `public`까지는 열지 않았다. + +### 5.3 `WebSocketConfig.java` — endpoint·broker·인터셉터 배선 + +```java +@EnableWebSocketMessageBroker +@Configuration +@RequiredArgsConstructor +public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { + + private final StompAuthChannelInterceptor stompAuthChannelInterceptor; + private final DashboardSubscriptionAuthorizationInterceptor dashboardSubscriptionAuthorizationInterceptor; + + @Override + public void registerStompEndpoints(StompEndpointRegistry registry) { + registry.addEndpoint("/ws") + .setAllowedOrigins(CorsConfig.ALLOWED_ORIGINS.toArray(new String[0])) + .withSockJS(); + } + + @Override + public void configureMessageBroker(MessageBrokerRegistry registry) { + registry.enableSimpleBroker("/topic"); + } + + @Override + public void configureClientInboundChannel(ChannelRegistration registration) { + registration.interceptors(stompAuthChannelInterceptor, dashboardSubscriptionAuthorizationInterceptor); + } +} +``` + +- `@EnableWebSocketMessageBroker` — STOMP 메시지 브로커 기능 자체를 켜는 스위치. 없으면 아래 + 콜백 메서드들이 호출되지 않는다. +- `registerStompEndpoints` — `/ws`로 연결 URL을 등록하고, 5.2에서 공유받은 origin 목록으로 + 제한한다. `.withSockJS()`는 순수 WebSocket이 막히는 환경(오래된 브라우저, 프록시)에서 HTTP + 폴링으로 자동 폴백해주는 옵션이다. +- `configureMessageBroker` — `/topic`으로 시작하는 목적지(`/topic/dashboard`)는 스프링 내장 + in-memory 브로커가 관리한다. 클라이언트가 구독해두면 서버가 `SimpMessagingTemplate`으로 보낸 + 메시지를 브로커가 구독자 전원에게 뿌린다(5.7). 클라이언트→서버 방향 메시지가 없어서 + `setApplicationDestinationPrefixes("/app")`는 넣지 않았다. +- `configureClientInboundChannel` — 클라이언트가 보내는 모든 STOMP 프레임이 지나가는 파이프에 + 두 인터셉터를 순서대로 건다. **순서가 중요하다** — CONNECT 인증(5.4)이 먼저 Principal을 + 세션에 부착해야, SUBSCRIBE 인가(5.5)가 그 Principal을 보고 판단할 수 있다. + +### 5.4 `StompAuthChannelInterceptor.java` — CONNECT 시점 JWT 인증 + +```java +@Component +@RequiredArgsConstructor +public class StompAuthChannelInterceptor implements ChannelInterceptor { + + private static final String AUTHORIZATION_HEADER = "Authorization"; + private static final String BEARER_PREFIX = "Bearer "; + + private final JwtProvider jwtProvider; + + @Override + @SuppressWarnings("unchecked") + public Message preSend(Message message, MessageChannel channel) { + StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); + + if (accessor != null && StompCommand.CONNECT.equals(accessor.getCommand())) { + String token = resolveToken(accessor); + Claims claims = token == null ? null : jwtProvider.getClaimsIfValid(token); + if (claims == null) { + throw new AccessDeniedException("유효하지 않은 인증 정보입니다."); + } + + String email = claims.getSubject(); + Long userId = claims.get("userId", Long.class); + List roles = (List) claims.get("roles"); + List authorities = roles.stream() + .map(role -> new SimpleGrantedAuthority("ROLE_" + role)) + .toList(); + + UsernamePasswordAuthenticationToken authentication = + new UsernamePasswordAuthenticationToken(email, null, authorities); + authentication.setDetails(userId); + accessor.setUser(authentication); + } + + return message; + } + + private String resolveToken(StompHeaderAccessor accessor) { + String bearer = accessor.getFirstNativeHeader(AUTHORIZATION_HEADER); + if (StringUtils.hasText(bearer) && bearer.startsWith(BEARER_PREFIX)) { + return bearer.substring(BEARER_PREFIX.length()); + } + return null; + } +} +``` + +WebSocket 자체는 그냥 양방향 파이프고 "이게 연결 요청인지 구독인지" 같은 구조가 없다. STOMP가 +그 위에 CONNECT·SUBSCRIBE·SEND 같은 프레임 타입을 정의해주는데, 클라이언트는 파이프가 열리면 +규격상 반드시 CONNECT를 제일 먼저 보내야 한다. 이 클래스는 그 첫 CONNECT 프레임을 가로채는 +문지기다. + +- `MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class)` — 이 메시지를 만들 + 때 `StompSubProtocolHandler`가 `leaveMutable(true)`로 캐시해 둔 **원본** Accessor를 가져온다. + `StompHeaderAccessor.wrap(message)`를 쓰면 검증 전용 **복사본**이 생겨서, 그 위에 `setUser()`를 + 호출해도 반환되는 `message`엔 반영되지 않고 조용히 사라진다 — 실제로 이렇게 짰다가 SUBSCRIBE + 단계에서 Principal이 매번 `null`로 나오는 버그를 냈다(4.3절). +- `resolveToken` — STOMP native header `Authorization: Bearer {token}`에서 토큰만 뽑는다. + `JwtAuthenticationFilter.resolveToken()`이 HTTP 헤더에서 하는 것과 동일한 규칙을, HTTP 헤더 + 대신 STOMP 프레임에서 하는 것만 다르다. +- `jwtProvider.getClaimsIfValid(token)` — 새 검증 로직이 아니라 기존 `JwtProvider`를 그대로 + 재사용한다. `null`이면(토큰 없음 또는 서명·만료 검증 실패) `AccessDeniedException`을 던져 + 연결을 즉시 끊는다. +- 유효하면 클레임에서 이메일·userId·역할을 꺼내 `Authentication`으로 조립하고 + `accessor.setUser(authentication)`으로 **세션에** 붙인다. HTTP는 요청 하나로 끝나 매번 + `SecurityContextHolder`를 새로 채우지만, WebSocket은 연결이 오래 유지되는 세션이라 이렇게 + 세션 레벨에 신원을 붙여두고 이후 모든 프레임에서 재사용한다. + +### 5.5 `DashboardSubscriptionAuthorizationInterceptor.java` — SUBSCRIBE 시점 ROLE_ADMIN 인가 + +```java +@Component +public class DashboardSubscriptionAuthorizationInterceptor implements ChannelInterceptor { + + private static final String DASHBOARD_TOPIC = "/topic/dashboard"; + private static final String ADMIN_AUTHORITY = "ROLE_ADMIN"; + + @Override + public Message preSend(Message message, MessageChannel channel) { + StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); + + if (accessor != null + && StompCommand.SUBSCRIBE.equals(accessor.getCommand()) + && DASHBOARD_TOPIC.equals(accessor.getDestination()) + && !isAdmin(accessor.getUser())) { + throw new AccessDeniedException("대시보드 구독 권한이 없습니다."); + } + + return message; + } + + private boolean isAdmin(Principal user) { + if (!(user instanceof Authentication authentication)) { + return false; + } + return authentication.getAuthorities().stream() + .map(GrantedAuthority::getAuthority) + .anyMatch(ADMIN_AUTHORITY::equals); + } +} +``` + +5.4가 "누구냐"를 확인했다면, 이 클래스는 "이 사람이 대시보드를 볼 자격이 있냐"만 본다. 5.4를 +통과해서 이미 세션이 열려있는(=로그인 확인된) 사람 중에서, `/topic/dashboard`를 구독하려 할 +때 ADMIN인지만 한 번 더 확인한다 — CONNECT 검증 하나에만 의존하지 않는 이중 방어다. + +`preSend`의 `if` 조건 4개가 전부 참이어야 거부한다: ① Accessor가 null이 아니고 ② SUBSCRIBE +프레임이고 ③ 목적지가 `/topic/dashboard`고 ④ `isAdmin()`이 `false`일 때. 넷 중 하나라도 아니면 +그냥 통과시킨다. + +`isAdmin`은 `accessor.getUser()`(타입은 `java.security.Principal`)를 `Authentication`으로 +다운캐스트해서 `getAuthorities()`를 봐야 권한 목록에 접근할 수 있다(`Principal` 자체는 이름만 +보장하고 권한 정보가 없다). `instanceof Authentication authentication` 패턴 매칭이 실패하면 +(이론상 CONNECT를 안 거친 경우) 그냥 "ADMIN 아님"으로 처리한다 — 애매하면 막는 fail-closed +방향이다. + +**원래 계획과 다르게 이 클래스가 존재하는 이유**는 4.2절 참고 — Spring Security의 +`@EnableWebSocketSecurity` DSL을 쓰려다 CSRF 문제로 포기하고 손으로 짠 결과물이다. + +### 5.6 `SecurityConfig.java` — `/ws/**`를 HTTP 레벨에서 열어둠 + +```diff + .requestMatchers("/auth/signup", "/auth/login").permitAll() ++ // WebSocket 핸드셰이크는 여기서 인증하지 않는다. 네이티브 websocket Transport는 ++ // Upgrade 요청에 커스텀 헤더를 실을 수 없어, 인증은 StompAuthChannelInterceptor가 ++ // STOMP CONNECT 프레임에서 담당하고 목적지별 인가는 DashboardSubscriptionAuthorizationInterceptor가 담당한다. ++ .requestMatchers("/ws/**").permitAll() + .requestMatchers("/admin/**").hasRole("ADMIN") + .anyRequest().authenticated() +``` + +`authorizeHttpRequests`는 위에서부터 순서대로 매칭되는 규칙 목록이다. `/ws/**`(WebSocket 연결을 +위한 HTTP Upgrade 요청)를 `.anyRequest().authenticated()`보다 **반드시 먼저** `permitAll()`로 +열어둬야 한다. 안 그러면 `Authorization` 헤더 없이 오는 최초 Upgrade 요청(네이티브 `websocket` +Transport의 경우 헤더를 아예 못 붙임)이 STOMP 프로토콜이 시작되기도 전에 401로 튕겨서, 5.4·5.5가 +아무리 잘 만들어져 있어도 도달할 기회 자체가 없다. + +이 앱은 `.csrf(AbstractHttpConfigurer::disable)` + `SessionCreationPolicy.STATELESS`라 세션도 +CSRF 토큰도 없는 상태인데, 이게 4.2절에서 설명하는 CSRF 버그의 근본 원인이기도 하다. + +### 5.7 `DashboardWebSocketController.java` — 실제 push 전송 + +```java +@Component +@RequiredArgsConstructor +public class DashboardWebSocketController { + + private static final String DASHBOARD_TOPIC = "/topic/dashboard"; + + private final SimpMessagingTemplate messagingTemplate; + + public void sendDashboardUpdate(DashboardSummaryResponse summary) { + messagingTemplate.convertAndSend(DASHBOARD_TOPIC, summary); + } +} +``` + +지금까지 5.3~5.6은 전부 "누가 연결·구독할 수 있는지"를 정하는 문지기였고, 이 클래스가 처음으로 +실제 데이터를 내보낸다. `convertAndSend("/topic/dashboard", summary)`를 호출하면 `summary`를 +JSON으로 직렬화해서 STOMP MESSAGE 프레임을 만들고, 5.3에서 켜둔 broker가 그 목적지를 구독 중인 +세션 전원(5.4·5.5를 통과한 ADMIN들)에게 뿌린다. + +`@RestController`가 아니라 `@Component`다 — HTTP로 직접 호출되는 엔드포인트가 아니라 다른 +서비스 코드가 자바 메서드처럼 호출하는 빈이다. 이름은 스펙 문서(F-OPS-02)의 명명을 그대로 +따랐다. + +이 메서드는 `DashboardQueryService`를 직접 호출해서 스스로 집계하지 않고, 최신 스냅샷을 파라미터로 +**받기만** 한다. 재처리 트리거·Worker 상태 전이 이벤트 등 호출하는 곳마다 "언제 다시 집계할지" +타이밍이 다르기 때문에, 이 클래스는 판단 없이 "받은 걸 보낸다"만 하는 얇은 전송 계층으로 남겨뒀다. + +**이번 이슈 시점에는 프로덕션 코드 어디에서도 이 메서드를 호출하지 않는다.** 통합 테스트(6.2절)가 +수동으로 호출해 push 자체가 동작하는지만 검증하고, 실제 호출 시점은 후속 이슈(재처리 트리거, +Worker 상태 전이 이벤트 훅)에서 채운다. + +## 6. 테스트 설계 + +### 6.1 단위 테스트 — `StompAuthChannelInterceptorTest` + +- 유효한 ADMIN 토큰 → CONNECT 프레임에 Principal 부착 +- `Authorization` 헤더 없음 → 거부 +- 토큰 무효 → 거부 +- CONNECT가 아닌 프레임 → 검증 없이 통과, `JwtProvider` 호출 안 됨 + +### 6.2 통합 테스트 — `DashboardWebSocketIntegrationTest` + +`@SpringBootTest(webEnvironment = RANDOM_PORT)` + 실제 `WebSocketStompClient`로 서버에 직접 +연결해서 검증한다. + +- ADMIN 토큰으로 구독 → `sendDashboardUpdate()` 호출 시 1초 내 수신 +- 토큰 없이 CONNECT → 연결 거부 (Transport 레벨 실패 또는 STOMP ERROR 프레임 중 어느 쪽이든 대응) +- USER(비 ADMIN) 토큰 → CONNECT는 성공하지만 `/topic/dashboard` SUBSCRIBE는 거부 + +## 7. 오류 계약 + +| 상황 | 처리 | +|---|---| +| `Authorization` 헤더 없이 CONNECT | `StompAuthChannelInterceptor`가 `AccessDeniedException` → 연결 종료 | +| 토큰 만료·서명 무효 | 위와 동일 | +| ADMIN이 아닌 사용자가 `/topic/dashboard` SUBSCRIBE | `DashboardSubscriptionAuthorizationInterceptor`가 `AccessDeniedException` → 구독 거부, 세션 종료 | +| `/ws/**` HTTP 핸드셰이크 자체 | 항상 permitAll, 여기서는 거부되지 않음 | + +## 8. 커밋 분할 + +1. `feat: #137 WebSocket 의존성 및 STOMP endpoint 설정 추가` +2. `feat: #137 STOMP CONNECT JWT 인증 Interceptor 추가` +3. `feat: #137 대시보드 SUBSCRIBE ROLE_ADMIN 인가 Interceptor 추가` +4. `feat: #137 대시보드 WebSocket push 전송 계층 추가` +5. `test: #137 WebSocket 인증·인가·push 단위·통합 테스트 추가` + +## 9. 완료 조건 + +- `/ws` 연결 및 `/topic/dashboard` 구독이 정상 동작한다 +- `sendDashboardUpdate()` 호출 시 구독 중인 클라이언트가 1초 내 최신 지표를 수신한다 +- 토큰 없음/무효 토큰으로 CONNECT 시 거부된다 +- ADMIN이 아닌 사용자는 `/topic/dashboard` SUBSCRIBE가 거부된다 +- 전체 빌드(`./gradlew build`)가 회귀 없이 통과한다 (733개 테스트, failures 0, errors 0) From ec4297c54d3e3c60297c3af9718c2c18f8794d9d Mon Sep 17 00:00:00 2001 From: kangcheolung Date: Mon, 10 Aug 2026 18:52:35 +0900 Subject: [PATCH 2/9] =?UTF-8?q?docs:=20#134=20=EC=84=A4=EA=B3=84=20?= =?UTF-8?q?=EB=AC=B8=EC=84=9C=EC=97=90=20=EC=8B=A4=EC=A0=9C=20=EC=BD=94?= =?UTF-8?q?=EB=93=9C=EC=99=80=20=EC=83=81=EC=84=B8=20=EC=84=A4=EB=AA=85=20?= =?UTF-8?q?=EB=B3=B4=EA=B0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Sonnet 5 --- ...gcheolung-#134-ragops-dashboard-summary.md | 284 ++++++++++++++++-- 1 file changed, 263 insertions(+), 21 deletions(-) diff --git a/docs/design/kangcheolung-#134-ragops-dashboard-summary.md b/docs/design/kangcheolung-#134-ragops-dashboard-summary.md index 161616c..e53660b 100644 --- a/docs/design/kangcheolung-#134-ragops-dashboard-summary.md +++ b/docs/design/kangcheolung-#134-ragops-dashboard-summary.md @@ -79,33 +79,275 @@ Authorization: Bearer {JWT} ## 4. 조회 구조 -### 4.1 Repository +읽는 순서는 Repository(4.1~4.3) → DTO(4.4) → Service(4.5) → Controller(4.6)를 따른다. +Service가 Repository들을 조합해서 DTO를 만들고, Controller는 그 결과를 그대로 HTTP로 감싸기만 +한다. -- `DocumentRepository`: `countByDeletedAtIsNull()`, `countByStatus(DocumentStatus)`, - `countByStatusIn(Collection)` — 메서드 이름 기반 자동 쿼리 -- `EmbeddingJobRepository`: `countByStatus(EmbeddingJobStatus)`, `findAllByStatus(EmbeddingJobStatus)` - (후속 전체 재처리 Issue에서 재사용 예정), `findAverageProcessingMillis()` — PostgreSQL - `EXTRACT(EPOCH FROM (completed_at - started_at)) * 1000` Native Query, 완료 Job이 없으면 `null` 반환 -- `SearchQueryRepository`: `countByCreatedAtAfter(LocalDateTime)` +### 4.1 `DocumentRepository` — 문서 카운트 3종 -### 4.2 DashboardQueryService +```java +// A담당자 영역 — B담당자는 존재 확인 등 읽기 전용으로만 사용 +public interface DocumentRepository extends JpaRepository { -`@Transactional(readOnly = true)` 클래스이며 자체 Repository 3개와 `WorkerNodeQueryService`에 -의존한다. 각 카테고리를 독립적으로 조회해 `DashboardSummaryResponse`로 조합한다. + /** + * 대시보드 집계 카드의 전체 문서 수. Soft-delete된 문서는 제외한다. + */ + long countByDeletedAtIsNull(); -- Worker 집계는 `resolveEffectiveStatus()`를 다시 구현하지 않고 `WorkerNodeQueryService.getWorkers()` - 호출 결과의 `status` 필드(이미 Heartbeat 기준으로 계산됨)를 그대로 센다. Heartbeat 판정 기준이 - 바뀌어도 이 Service는 수정할 필요가 없다. -- `avgProcessMs`는 Native Query가 `null`을 반환하면 그대로 `null`을 응답하고, 값이 있으면 반올림해 - `Long`으로 변환한다. 0으로 기본값을 채우지 않는다 — 완료 Job이 없는 상태에서 "평균 0ms"는 사실과 - 다른 정보이기 때문이다. -- 최근 24시간 기준 시각은 주입받은 `Clock`으로 계산해 테스트 시 고정 가능하게 한다. + /** + * 대시보드 집계 카드에서 특정 상태 하나에 속하는 문서 수를 센다 (예: 검색 가능 문서 수). + */ + long countByStatus(DocumentStatus status); -### 4.3 DTO + /** + * 대시보드 집계 카드에서 여러 상태에 걸친 문서 수를 센다 (예: 인덱싱 대기 중 문서 수). + */ + long countByStatusIn(Collection statuses); -`DashboardSummaryResponse`가 `DocumentsSummaryResponse` / `JobsSummaryResponse` / -`WorkersSummaryResponse` / `SearchSummaryResponse` 4개를 필드로 갖는다. 각 필드는 `@Schema`로 -Swagger 설명을 붙인다. + // ... 기존 메서드들 +} +``` + +세 메서드 다 `@Query` 없이 메서드 이름만으로 Spring Data JPA가 쿼리를 자동 생성한다 +(`countByDeletedAtIsNull` → `WHERE deleted_at IS NULL`, `countByStatusIn` → `WHERE status IN (...)`). +`DocumentRepository`는 파일 상단 주석에 "A담당자 영역 — B담당자는 읽기 전용으로만 사용"이라고 +이미 명시돼 있어서, 쓰기 메서드는 추가하지 않고 count류만 붙였다. + +### 4.2 `EmbeddingJobRepository` — 상태별 카운트 + 평균 처리 시간 + +```java +/** + * 대시보드 집계 카드(대기/처리 중/실패 작업 수)에 사용하는 상태별 Job 수를 센다. + */ +long countByStatus(EmbeddingJobStatus status); + +/** + * 관리자 전체 재처리 대상인 FAILED Job 전체를 조회한다. + */ +List findAllByStatus(EmbeddingJobStatus status); + +/** + * 완료된 Job의 평균 처리 시간을 밀리초 단위로 계산한다. + * + *

Queue 대기 시간({@code created_at})은 제외하고 Worker가 실제로 처리한 구간({@code started_at} + * ~ {@code completed_at})만 반영한다. 완료된 Job이 없으면 {@code null}을 반환한다. + */ +@Query( + value = """ + SELECT AVG(EXTRACT(EPOCH FROM (completed_at - started_at)) * 1000) + FROM embedding_jobs + WHERE status = 'INDEXED' + AND started_at IS NOT NULL + AND completed_at IS NOT NULL + """, + nativeQuery = true +) +Double findAverageProcessingMillis(); +``` + +`countByStatus`는 4.1과 같은 메서드 이름 자동 쿼리. `findAllByStatus`는 이번 이슈에서 직접 +쓰지 않고, 후속 이슈(FAILED 전체 재처리)가 `EmbeddingJobRepository.findAllByStatus(FAILED)`로 +재사용할 걸 미리 준비해둔 것이다. + +`findAverageProcessingMillis()`는 `AVG(completed_at - started_at)`을 쓰지 `AVG(completed_at - +created_at)`을 쓰지 않는다 — `created_at`부터 재면 Queue에서 대기한 시간까지 "처리 시간"에 섞여서 +지표 의미가 흐려진다. PostgreSQL 전용 함수 `EXTRACT(EPOCH FROM ...)`를 쓰기 때문에 JPQL이 아니라 +`nativeQuery = true`로 짰다. `AVG`는 대상 행이 0개면 SQL 표준상 `NULL`을 반환하므로, 반환 타입도 +기본값 `0.0`이 아니라 `Double`(nullable)로 선언해서 "완료된 Job이 아예 없다"는 사실을 그대로 +드러낸다. + +### 4.3 `SearchQueryRepository` — 최근 검색 수 + +```java +public interface SearchQueryRepository extends JpaRepository { + + /** + * 대시보드 집계 카드의 최근 검색 요청 수. 기준 시각 이후 생성된 검색 Query를 센다. + */ + long countByCreatedAtAfter(LocalDateTime since); +} +``` + +원래 이 인터페이스는 `JpaRepository`만 상속하고 메서드가 하나도 없었다. 이번에 처음 추가한 +메서드다. `since` 기준 시각(now - 24h)은 4.5절 Service가 계산해서 넘긴다. + +### 4.4 DTO — `DashboardSummaryResponse` + 하위 4개 + +```java +public record DashboardSummaryResponse( + @Schema(description = "문서 현황") + DocumentsSummaryResponse documents, + + @Schema(description = "인덱싱 작업 현황") + JobsSummaryResponse jobs, + + @Schema(description = "Worker 현황") + WorkersSummaryResponse workers, + + @Schema(description = "검색 현황") + SearchSummaryResponse search +) { } + +public record DocumentsSummaryResponse( + @Schema(description = "전체 문서 수 (Soft-delete 제외)", example = "25368") + long total, + + @Schema(description = "검색 가능 문서 수 (INDEXED 상태)", example = "21742") + long searchable, + + @Schema(description = "인덱싱 대기 중인 문서 수 (UPLOADED, INDEXING 상태)", example = "132") + long pendingIndex +) { } + +public record JobsSummaryResponse( + @Schema(description = "인덱싱 대기 작업 수 (PENDING)", example = "132") + long pending, + + @Schema(description = "처리 중인 작업 수 (PROCESSING)", example = "8") + long processing, + + @Schema(description = "실패 작업 수 (FAILED)", example = "27") + long failed, + + @Schema(description = "평균 임베딩 처리 시간(ms). Queue 대기 시간은 제외한 순수 처리 시간이며, " + + "완료된 Job이 없으면 null", example = "3200") + Long avgProcessMs +) { } + +public record WorkersSummaryResponse( + @Schema(description = "정상(ACTIVE·IDLE) Worker 수", example = "5") + long activeCount, + + @Schema(description = "전체 등록 Worker 수", example = "6") + long totalCount +) { } + +public record SearchSummaryResponse( + @Schema(description = "최근 24시간 검색 요청 수", example = "342") + long recent24hCount +) { } +``` + +5개 파일로 나눈 이유는 프로젝트 컨벤션(`dto/response`에 응답 DTO 하나당 파일 하나)을 따른 것이다. +`documents`/`jobs`/`workers`/`search`가 각자 독립된 record라, 나중에 특정 카테고리 하나만 API로 +따로 빼야 할 일이 생겨도 재사용하기 쉽다. `avgProcessMs`만 `Long`(기본 타입 `long`이 아니라 +Wrapper)인 이유는 4.2에서 설명한 `null` 가능성을 DTO까지 그대로 전달하기 위해서다. + +### 4.5 `DashboardQueryService` + +```java +@Transactional(readOnly = true) +@Service +@RequiredArgsConstructor +public class DashboardQueryService { + + private static final List PENDING_INDEX_STATUSES = + List.of(DocumentStatus.UPLOADED, DocumentStatus.INDEXING); + + private final DocumentRepository documentRepository; + private final EmbeddingJobRepository embeddingJobRepository; + private final SearchQueryRepository searchQueryRepository; + private final WorkerNodeQueryService workerNodeQueryService; + private final Clock clock; + + public DashboardSummaryResponse getSummary() { + return new DashboardSummaryResponse( + getDocumentsSummary(), + getJobsSummary(), + getWorkersSummary(), + getSearchSummary() + ); + } + + private DocumentsSummaryResponse getDocumentsSummary() { + return new DocumentsSummaryResponse( + documentRepository.countByDeletedAtIsNull(), + documentRepository.countByStatus(DocumentStatus.INDEXED), + documentRepository.countByStatusIn(PENDING_INDEX_STATUSES) + ); + } + + private JobsSummaryResponse getJobsSummary() { + Double averageMillis = embeddingJobRepository.findAverageProcessingMillis(); + return new JobsSummaryResponse( + embeddingJobRepository.countByStatus(EmbeddingJobStatus.PENDING), + embeddingJobRepository.countByStatus(EmbeddingJobStatus.PROCESSING), + embeddingJobRepository.countByStatus(EmbeddingJobStatus.FAILED), + averageMillis == null ? null : Math.round(averageMillis) + ); + } + + private WorkersSummaryResponse getWorkersSummary() { + List workers = workerNodeQueryService.getWorkers(); + long activeCount = workers.stream() + .filter(worker -> worker.status() == WorkerStatus.ACTIVE || worker.status() == WorkerStatus.IDLE) + .count(); + return new WorkersSummaryResponse(activeCount, workers.size()); + } + + private SearchSummaryResponse getSearchSummary() { + LocalDateTime since = LocalDateTime.now(clock).minusHours(24); + return new SearchSummaryResponse(searchQueryRepository.countByCreatedAtAfter(since)); + } +} +``` + +- `getSummary()`가 4개의 `private` 메서드를 호출해서 `DashboardSummaryResponse`를 조립하는 + 단순 오케스트레이션이다. 카테고리 4개가 서로 의존하지 않아 순서는 상관없다. +- `getDocumentsSummary()`는 4.1의 세 메서드를 그대로 호출한다. `PENDING_INDEX_STATUSES`를 + 상수로 뺀 건 `pendingIndex`의 정의(`UPLOADED`, `INDEXING`)가 이 클래스 밖에서도 참조될 일이 + 없어서 필드 상수로 충분하다고 판단했다. +- `getJobsSummary()`가 `averageMillis == null ? null : Math.round(averageMillis)`로 널 처리를 + 명시적으로 한다 — `Math.round(null)`은 컴파일이 안 되고, 여기서 `0`으로 기본값을 채우면 "완료 + Job이 없다"와 "평균이 정확히 0ms다"를 구분 못 하게 된다. +- `getWorkersSummary()`가 이 Service에서 유일하게 자기 Repository가 아니라 **다른 도메인의 + Query Service**(`WorkerNodeQueryService`)를 부른다. `WorkerNodeQueryService.getWorkers()`가 + 반환하는 `WorkerNodeResponse.status()`는 이미 Heartbeat 기준으로 계산된 값이라(`WorkerNode.resolveEffectiveStatus()`), + 여기서 그 판정 로직을 다시 만들 필요가 없다 — `ACTIVE`/`IDLE`인 것만 세면 된다. +- `getSearchSummary()`는 `LocalDateTime.now()`를 직접 안 쓰고 주입받은 `Clock`으로 현재 시각을 + 구한다. 이러면 테스트에서 `Clock.fixed(...)`로 시간을 고정해 "24시간 전"을 결정론적으로 + 검증할 수 있다(6.1절 테스트 참고). + +### 4.6 `DashboardController` + +```java +@Tag(name = "Admin - Dashboard", description = "관리자 전용 RAGOps Dashboard 집계 지표 API") +@RestController +@RequestMapping("/admin/dashboard") +@RequiredArgsConstructor +public class DashboardController { + + private final DashboardQueryService dashboardQueryService; + + @Operation( + summary = "대시보드 집계 지표 조회", + description = "문서·인덱싱 작업·Worker·검색 현황을 하나의 응답으로 집계해서 반환합니다. " + + "모든 지표는 조회 시점 기준 Snapshot이며, 실시간 WebSocket push는 이 API의 범위가 아닙니다." + ) + @ApiResponses({ + @ApiResponse(responseCode = "200", description = "집계 지표 조회 성공"), + @ApiResponse( + responseCode = "403", + description = "인증되지 않았거나 ADMIN 권한 없음", + content = @Content(schema = @Schema(implementation = ErrorResponse.class)) + ) + }) + @GetMapping(value = "/summary", produces = MediaType.APPLICATION_JSON_VALUE) + public ResponseEntity> getSummary() { + return ResponseUtils.ok(dashboardQueryService.getSummary()); + } +} +``` + +`@RequestMapping("/admin/dashboard")` + `@GetMapping("/summary")`로 최종 경로가 +`/admin/dashboard/summary`가 된다. `/admin/**`는 `SecurityConfig`에 이미 `hasRole("ADMIN")` +규칙이 있어서 이 Controller엔 별도 권한 어노테이션이 없다 — 경로만 `/admin` 아래 두면 자동으로 +ADMIN 전용이 된다. 메서드 본문은 `dashboardQueryService.getSummary()` 결과를 +`ResponseUtils.ok()`로 감싸는 게 전부다 — 입력 파라미터가 없어서 검증 로직도 없다. + +`description`의 마지막 문장("실시간 WebSocket push는 이 API의 범위가 아닙니다")은 CodeRabbit +리뷰로 정정한 부분이다. 원래는 "WebSocket push로 제공됩니다"라고 써서, 아직 구현되지도 않은 +기능(이슈#137에서 구현)이 이미 있는 것처럼 문서화하는 실수가 있었다. ## 5. 오류 계약 From 6dacb1ed9de91a620e69338c51a26ffe8ffb664d Mon Sep 17 00:00:00 2001 From: kangcheolung Date: Mon, 10 Aug 2026 18:52:50 +0900 Subject: [PATCH 3/9] =?UTF-8?q?feat:=20#137=20WebSocket=20=EC=9D=98?= =?UTF-8?q?=EC=A1=B4=EC=84=B1=20=EB=B0=8F=20CORS=20origin=20=EA=B3=B5?= =?UTF-8?q?=EC=9C=A0=20=EC=84=A4=EC=A0=95=20=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Sonnet 5 --- build.gradle | 1 + .../java/com/opensource/docgrid/global/config/CorsConfig.java | 3 ++- 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/build.gradle b/build.gradle index 71c2236..6a21288 100644 --- a/build.gradle +++ b/build.gradle @@ -28,6 +28,7 @@ dependencies { implementation 'org.springframework.boot:spring-boot-starter-security' implementation 'org.springframework.boot:spring-boot-starter-validation' implementation 'org.springframework.boot:spring-boot-starter-web' + implementation 'org.springframework.boot:spring-boot-starter-websocket' implementation 'org.springframework.ai:spring-ai-starter-mcp-server-webmvc' implementation 'org.springdoc:springdoc-openapi-starter-webmvc-ui:2.8.9' implementation 'io.minio:minio:8.5.17' diff --git a/src/main/java/com/opensource/docgrid/global/config/CorsConfig.java b/src/main/java/com/opensource/docgrid/global/config/CorsConfig.java index 06abe61..9d395bd 100644 --- a/src/main/java/com/opensource/docgrid/global/config/CorsConfig.java +++ b/src/main/java/com/opensource/docgrid/global/config/CorsConfig.java @@ -13,7 +13,8 @@ @Configuration public class CorsConfig implements WebMvcConfigurer { - private static final List ALLOWED_ORIGINS = List.of( + // WebSocketConfig가 STOMP endpoint 허용 origin으로 재사용하므로 package-private으로 둔다. + static final List ALLOWED_ORIGINS = List.of( "http://localhost:3000", "http://localhost:8080" ); From 83b6b8443a7d676ae8e606563871e67093374c6e Mon Sep 17 00:00:00 2001 From: kangcheolung Date: Mon, 10 Aug 2026 18:53:02 +0900 Subject: [PATCH 4/9] =?UTF-8?q?feat:=20#137=20STOMP=20CONNECT=20JWT=20?= =?UTF-8?q?=EC=9D=B8=EC=A6=9D=20Interceptor=20=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Sonnet 5 --- .../auth/jwt/StompAuthChannelInterceptor.java | 92 +++++++++++++++++++ 1 file changed, 92 insertions(+) create mode 100644 src/main/java/com/opensource/docgrid/domain/auth/jwt/StompAuthChannelInterceptor.java diff --git a/src/main/java/com/opensource/docgrid/domain/auth/jwt/StompAuthChannelInterceptor.java b/src/main/java/com/opensource/docgrid/domain/auth/jwt/StompAuthChannelInterceptor.java new file mode 100644 index 0000000..e373b4a --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/auth/jwt/StompAuthChannelInterceptor.java @@ -0,0 +1,92 @@ +package com.opensource.docgrid.domain.auth.jwt; + +import java.util.List; + +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.simp.stomp.StompCommand; +import org.springframework.messaging.simp.stomp.StompHeaderAccessor; +import org.springframework.messaging.support.ChannelInterceptor; +import org.springframework.messaging.support.MessageHeaderAccessor; +import org.springframework.security.access.AccessDeniedException; +import org.springframework.security.authentication.UsernamePasswordAuthenticationToken; +import org.springframework.security.core.authority.SimpleGrantedAuthority; +import org.springframework.stereotype.Component; +import org.springframework.util.StringUtils; + +import io.jsonwebtoken.Claims; +import lombok.RequiredArgsConstructor; + +/** + * STOMP CONNECT 프레임의 JWT를 검증해 WebSocket 세션에 Principal을 부착한다. + * + *

WebSocket 자체는 그냥 양방향 파이프를 열어줄 뿐, "이 메시지가 연결 요청인지 구독인지" 같은 + * 구조가 없다. STOMP는 그 파이프 위에 CONNECT·SUBSCRIBE·SEND 같은 프레임 타입을 정의해 의미 + * 있는 대화를 가능하게 하는 프로토콜이고, 클라이언트는 파이프가 열리면 규격상 반드시 CONNECT + * 프레임을 제일 먼저 보내야 한다. 이 Interceptor는 그 첫 CONNECT 프레임을 가로채 "로그인한 + * 사용자인가"를 확인하는 문지기이며, JWT가 없거나 무효하면 그 자리에서 연결을 끊는다. + * + *

HTTP 핸드셰이크(/ws)는 {@code SecurityConfig}에서 permitAll로 열려 있다. 네이티브 + * {@code websocket} Transport는 Upgrade 요청에 커스텀 헤더를 실을 수 없어 HTTP 레벨 인증이 + * Transport 종류에 따라 되다 안되다 하므로, 인증은 Transport와 무관하게 항상 커스텀 헤더를 실을 + * 수 있는 STOMP CONNECT 프레임으로 옮긴다. SUBSCRIBE 권한 검증은 {@link + * com.opensource.docgrid.domain.dashboard.websocket.DashboardSubscriptionAuthorizationInterceptor}가 + * 별도로 담당한다. + * + *

토큰 파싱과 {@code JwtProvider} 검증은 {@code JwtAuthenticationFilter}의 HTTP 경로와 + * 같은 규칙을 그대로 재사용한다. 다만 결과를 담는 곳이 다르다 — HTTP는 요청 하나로 끝나 매번 + * {@code SecurityContextHolder}를 새로 채우지만, WebSocket은 연결이 오래 유지되는 세션이라 + * {@code accessor.setUser()}로 세션 자체에 Principal을 붙여 이후 모든 프레임에서 재사용한다. + * + *

이때 Accessor는 반드시 {@link MessageHeaderAccessor#getAccessor}로 가져와야 한다. + * {@code StompHeaderAccessor.wrap(message)}는 검증 전용 복사본이라 그 위에 {@code setUser()}를 + * 호출해도 반환되는 {@code message}에는 반영되지 않고, SUBSCRIBE 단계에서 Principal이 조용히 + * 사라지는 문제로 이어진다. + */ +@Component +@RequiredArgsConstructor +public class StompAuthChannelInterceptor implements ChannelInterceptor { + + private static final String AUTHORIZATION_HEADER = "Authorization"; + private static final String BEARER_PREFIX = "Bearer "; + + private final JwtProvider jwtProvider; + + @Override + @SuppressWarnings("unchecked") + public Message preSend(Message message, MessageChannel channel) { + StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); + + // 인증은 세션 시작 시점(CONNECT) 한 번만 하면 된다. SUBSCRIBE 등 이후 프레임은 그대로 통과시킨다. + if (accessor != null && StompCommand.CONNECT.equals(accessor.getCommand())) { + String token = resolveToken(accessor); + // JWT가 유효해야 신원을 확인한 것으로 본다. 없거나 무효하면 여기서 바로 연결을 끊는다. + Claims claims = token == null ? null : jwtProvider.getClaimsIfValid(token); + if (claims == null) { + throw new AccessDeniedException("유효하지 않은 인증 정보입니다."); + } + + String email = claims.getSubject(); + Long userId = claims.get("userId", Long.class); + List roles = (List) claims.get("roles"); + List authorities = roles.stream() + .map(role -> new SimpleGrantedAuthority("ROLE_" + role)) + .toList(); + + UsernamePasswordAuthenticationToken authentication = + new UsernamePasswordAuthenticationToken(email, null, authorities); + authentication.setDetails(userId); + accessor.setUser(authentication); + } + + return message; + } + + private String resolveToken(StompHeaderAccessor accessor) { + String bearer = accessor.getFirstNativeHeader(AUTHORIZATION_HEADER); + if (StringUtils.hasText(bearer) && bearer.startsWith(BEARER_PREFIX)) { + return bearer.substring(BEARER_PREFIX.length()); + } + return null; + } +} From d13b3947c46a74273467f265f7e0cc7a1ea0148c Mon Sep 17 00:00:00 2001 From: kangcheolung Date: Mon, 10 Aug 2026 18:53:13 +0900 Subject: [PATCH 5/9] =?UTF-8?q?feat:=20#137=20=EB=8C=80=EC=8B=9C=EB=B3=B4?= =?UTF-8?q?=EB=93=9C=20SUBSCRIBE=20ROLE=5FADMIN=20=EC=9D=B8=EA=B0=80=20Int?= =?UTF-8?q?erceptor=20=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Sonnet 5 --- ...dSubscriptionAuthorizationInterceptor.java | 55 +++++++++++++++++++ 1 file changed, 55 insertions(+) create mode 100644 src/main/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardSubscriptionAuthorizationInterceptor.java diff --git a/src/main/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardSubscriptionAuthorizationInterceptor.java b/src/main/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardSubscriptionAuthorizationInterceptor.java new file mode 100644 index 0000000..3930e10 --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardSubscriptionAuthorizationInterceptor.java @@ -0,0 +1,55 @@ +package com.opensource.docgrid.domain.dashboard.websocket; + +import java.security.Principal; + +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.simp.stomp.StompCommand; +import org.springframework.messaging.simp.stomp.StompHeaderAccessor; +import org.springframework.messaging.support.ChannelInterceptor; +import org.springframework.messaging.support.MessageHeaderAccessor; +import org.springframework.security.access.AccessDeniedException; +import org.springframework.security.core.Authentication; +import org.springframework.security.core.GrantedAuthority; +import org.springframework.stereotype.Component; + +/** + * {@code /topic/dashboard} SUBSCRIBE 요청에 ROLE_ADMIN 권한을 요구한다. + * + *

{@code StompAuthChannelInterceptor}가 CONNECT 시점에 세션에 부착한 Principal을 재사용해 + * 목적지 접근 시점에 다시 한 번 검증한다. CONNECT 검증 하나에만 의존하지 않는 이중 방어다. + * + *

{@code @EnableWebSocketSecurity}(Spring Security 메시지 인가 DSL)는 STOMP endpoint가 + * 등록된 것을 감지하면 세션 기반 CSRF 토큰을 무조건 요구하는 {@code CsrfChannelInterceptor}를 + * 함께 붙인다. 이 앱은 세션이 없는 stateless JWT 인증이라 CSRF 토큰이 존재할 수 없어 모든 CONNECT가 + * {@code MissingCsrfTokenException}으로 거부되므로, 그 DSL 대신 이 수동 Interceptor로 구현한다. + */ +@Component +public class DashboardSubscriptionAuthorizationInterceptor implements ChannelInterceptor { + + private static final String DASHBOARD_TOPIC = "/topic/dashboard"; + private static final String ADMIN_AUTHORITY = "ROLE_ADMIN"; + + @Override + public Message preSend(Message message, MessageChannel channel) { + StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); + + if (accessor != null + && StompCommand.SUBSCRIBE.equals(accessor.getCommand()) + && DASHBOARD_TOPIC.equals(accessor.getDestination()) + && !isAdmin(accessor.getUser())) { + throw new AccessDeniedException("대시보드 구독 권한이 없습니다."); + } + + return message; + } + + private boolean isAdmin(Principal user) { + if (!(user instanceof Authentication authentication)) { + return false; + } + return authentication.getAuthorities().stream() + .map(GrantedAuthority::getAuthority) + .anyMatch(ADMIN_AUTHORITY::equals); + } +} From 34936182cdf6a45dfbe83f0362ce896dbdb74c3b Mon Sep 17 00:00:00 2001 From: kangcheolung Date: Mon, 10 Aug 2026 18:53:27 +0900 Subject: [PATCH 6/9] =?UTF-8?q?feat:=20#137=20STOMP=20endpoint=20=EC=84=A4?= =?UTF-8?q?=EC=A0=95=20=EB=B0=8F=20HTTP=20=EB=A0=88=EB=B2=A8=20=EC=9D=B8?= =?UTF-8?q?=EC=A6=9D=20=EB=B0=B0=EC=84=A0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Sonnet 5 --- .../docgrid/global/config/SecurityConfig.java | 4 ++ .../global/config/WebSocketConfig.java | 47 +++++++++++++++++++ 2 files changed, 51 insertions(+) create mode 100644 src/main/java/com/opensource/docgrid/global/config/WebSocketConfig.java diff --git a/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java b/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java index 1aeab5d..5c5032a 100644 --- a/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java +++ b/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java @@ -39,6 +39,10 @@ public SecurityFilterChain filterChain(HttpSecurity http) throws Exception { .requestMatchers("/swagger-ui/**", "/v3/api-docs/**").permitAll() .requestMatchers("/departments").permitAll() .requestMatchers("/auth/signup", "/auth/login").permitAll() + // WebSocket 핸드셰이크는 여기서 인증하지 않는다. 네이티브 websocket Transport는 + // Upgrade 요청에 커스텀 헤더를 실을 수 없어, 인증은 StompAuthChannelInterceptor가 + // STOMP CONNECT 프레임에서 담당하고 목적지별 인가는 DashboardSubscriptionAuthorizationInterceptor가 담당한다. + .requestMatchers("/ws/**").permitAll() .requestMatchers("/admin/**").hasRole("ADMIN") .anyRequest().authenticated() ) diff --git a/src/main/java/com/opensource/docgrid/global/config/WebSocketConfig.java b/src/main/java/com/opensource/docgrid/global/config/WebSocketConfig.java new file mode 100644 index 0000000..619dd04 --- /dev/null +++ b/src/main/java/com/opensource/docgrid/global/config/WebSocketConfig.java @@ -0,0 +1,47 @@ +package com.opensource.docgrid.global.config; + +import org.springframework.context.annotation.Configuration; +import org.springframework.messaging.simp.config.ChannelRegistration; +import org.springframework.messaging.simp.config.MessageBrokerRegistry; +import org.springframework.web.socket.config.annotation.EnableWebSocketMessageBroker; +import org.springframework.web.socket.config.annotation.StompEndpointRegistry; +import org.springframework.web.socket.config.annotation.WebSocketMessageBrokerConfigurer; + +import com.opensource.docgrid.domain.auth.jwt.StompAuthChannelInterceptor; +import com.opensource.docgrid.domain.dashboard.websocket.DashboardSubscriptionAuthorizationInterceptor; + +import lombok.RequiredArgsConstructor; + +/** + * RAGOps Dashboard 실시간 push를 위한 STOMP endpoint와 Message Broker 설정. + * + *

인증·인가는 이 설정이 아니라 {@link StompAuthChannelInterceptor}(CONNECT 시점 인증)와 + * {@link DashboardSubscriptionAuthorizationInterceptor}(SUBSCRIBE 시점 인가)가 담당한다. + * 이 클래스는 전송 계층 구성(endpoint·broker·origin)과 두 Interceptor의 등록 순서만 책임진다. + */ +@EnableWebSocketMessageBroker +@Configuration +@RequiredArgsConstructor +public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { + + private final StompAuthChannelInterceptor stompAuthChannelInterceptor; + private final DashboardSubscriptionAuthorizationInterceptor dashboardSubscriptionAuthorizationInterceptor; + + @Override + public void registerStompEndpoints(StompEndpointRegistry registry) { + registry.addEndpoint("/ws") + .setAllowedOrigins(CorsConfig.ALLOWED_ORIGINS.toArray(new String[0])) + .withSockJS(); + } + + @Override + public void configureMessageBroker(MessageBrokerRegistry registry) { + registry.enableSimpleBroker("/topic"); + } + + @Override + public void configureClientInboundChannel(ChannelRegistration registration) { + // CONNECT 인증이 SUBSCRIBE 인가보다 먼저 Principal을 세션에 부착해야 하므로 순서를 고정한다. + registration.interceptors(stompAuthChannelInterceptor, dashboardSubscriptionAuthorizationInterceptor); + } +} From 38073704d419ea260adf568924d9112ec31e6520 Mon Sep 17 00:00:00 2001 From: kangcheolung Date: Mon, 10 Aug 2026 18:53:41 +0900 Subject: [PATCH 7/9] =?UTF-8?q?feat:=20#137=20=EB=8C=80=EC=8B=9C=EB=B3=B4?= =?UTF-8?q?=EB=93=9C=20WebSocket=20push=20=EC=A0=84=EC=86=A1=20=EA=B3=84?= =?UTF-8?q?=EC=B8=B5=20=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Sonnet 5 --- .../DashboardWebSocketController.java | 28 +++++++++++++++++++ 1 file changed, 28 insertions(+) create mode 100644 src/main/java/com/opensource/docgrid/domain/dashboard/controller/DashboardWebSocketController.java diff --git a/src/main/java/com/opensource/docgrid/domain/dashboard/controller/DashboardWebSocketController.java b/src/main/java/com/opensource/docgrid/domain/dashboard/controller/DashboardWebSocketController.java new file mode 100644 index 0000000..b3ee6c0 --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/dashboard/controller/DashboardWebSocketController.java @@ -0,0 +1,28 @@ +package com.opensource.docgrid.domain.dashboard.controller; + +import org.springframework.messaging.simp.SimpMessagingTemplate; +import org.springframework.stereotype.Component; + +import com.opensource.docgrid.domain.dashboard.dto.response.DashboardSummaryResponse; + +import lombok.RequiredArgsConstructor; + +/** + * RAGOps Dashboard 집계 지표를 {@code /topic/dashboard} 구독자에게 push하는 전송 계층. + * + *

이 컨트롤러는 직접 지표를 집계하지 않는다. 호출하는 쪽(재처리 트리거, 상태 전이 이벤트 + * 리스너 등)이 {@code DashboardQueryService}로 최신 Snapshot을 계산해 넘겨주면 그대로 + * 브로드캐스트만 한다. + */ +@Component +@RequiredArgsConstructor +public class DashboardWebSocketController { + + private static final String DASHBOARD_TOPIC = "/topic/dashboard"; + + private final SimpMessagingTemplate messagingTemplate; + + public void sendDashboardUpdate(DashboardSummaryResponse summary) { + messagingTemplate.convertAndSend(DASHBOARD_TOPIC, summary); + } +} From b93bf9116897bcf942a4295280c81b4bce81d6db Mon Sep 17 00:00:00 2001 From: kangcheolung Date: Mon, 10 Aug 2026 18:53:54 +0900 Subject: [PATCH 8/9] =?UTF-8?q?test:=20#137=20WebSocket=20=EC=9D=B8?= =?UTF-8?q?=EC=A6=9D=C2=B7=EC=9D=B8=EA=B0=80=C2=B7push=20=EB=8B=A8?= =?UTF-8?q?=EC=9C=84=C2=B7=ED=86=B5=ED=95=A9=20=ED=85=8C=EC=8A=A4=ED=8A=B8?= =?UTF-8?q?=20=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Sonnet 5 --- .../jwt/StompAuthChannelInterceptorTest.java | 116 ++++++++++ .../DashboardWebSocketIntegrationTest.java | 210 ++++++++++++++++++ 2 files changed, 326 insertions(+) create mode 100644 src/test/java/com/opensource/docgrid/domain/auth/jwt/StompAuthChannelInterceptorTest.java create mode 100644 src/test/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardWebSocketIntegrationTest.java diff --git a/src/test/java/com/opensource/docgrid/domain/auth/jwt/StompAuthChannelInterceptorTest.java b/src/test/java/com/opensource/docgrid/domain/auth/jwt/StompAuthChannelInterceptorTest.java new file mode 100644 index 0000000..2d1aa96 --- /dev/null +++ b/src/test/java/com/opensource/docgrid/domain/auth/jwt/StompAuthChannelInterceptorTest.java @@ -0,0 +1,116 @@ +package com.opensource.docgrid.domain.auth.jwt; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.mock; +import static org.mockito.BDDMockito.then; + +import java.security.Principal; +import java.util.List; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.simp.stomp.StompCommand; +import org.springframework.messaging.simp.stomp.StompHeaderAccessor; +import org.springframework.messaging.support.MessageBuilder; +import org.springframework.security.access.AccessDeniedException; +import org.springframework.security.core.Authentication; +import org.springframework.security.core.GrantedAuthority; + +import io.jsonwebtoken.Claims; + +@ExtendWith(MockitoExtension.class) +@DisplayName("StompAuthChannelInterceptor 단위 테스트") +class StompAuthChannelInterceptorTest { + + @Mock private JwtProvider jwtProvider; + @Mock private MessageChannel channel; + + private StompAuthChannelInterceptor interceptor; + + @BeforeEach + void setUp() { + interceptor = new StompAuthChannelInterceptor(jwtProvider); + } + + @Test + @DisplayName("정상 케이스: 유효한 ADMIN 토큰이면 CONNECT 프레임에 Principal을 부착한다") + void preSend_attachesPrincipal_whenTokenValid() { + // Given + Claims claims = mock(Claims.class); + given(claims.getSubject()).willReturn("admin@example.com"); + given(claims.get("userId", Long.class)).willReturn(1L); + given(claims.get("roles")).willReturn(List.of("ADMIN")); + given(jwtProvider.getClaimsIfValid("valid-token")).willReturn(claims); + + Message connectMessage = connectMessage("Bearer valid-token"); + + // When + Message result = interceptor.preSend(connectMessage, channel); + + // Then + StompHeaderAccessor resultAccessor = StompHeaderAccessor.wrap(result); + Principal user = resultAccessor.getUser(); + assertThat(user).isNotNull(); + assertThat(user.getName()).isEqualTo("admin@example.com"); + assertThat(((Authentication) user).getAuthorities()) + .extracting(GrantedAuthority::getAuthority) + .containsExactly("ROLE_ADMIN"); + } + + @Test + @DisplayName("예외 케이스: Authorization 헤더가 없으면 연결을 거부한다") + void preSend_throws_whenNoAuthorizationHeader() { + // Given + Message connectMessage = connectMessage(null); + + // When & Then + assertThatThrownBy(() -> interceptor.preSend(connectMessage, channel)) + .isInstanceOf(AccessDeniedException.class); + then(jwtProvider).shouldHaveNoInteractions(); + } + + @Test + @DisplayName("예외 케이스: 토큰이 유효하지 않으면 연결을 거부한다") + void preSend_throws_whenTokenInvalid() { + // Given + given(jwtProvider.getClaimsIfValid("invalid-token")).willReturn(null); + Message connectMessage = connectMessage("Bearer invalid-token"); + + // When & Then + assertThatThrownBy(() -> interceptor.preSend(connectMessage, channel)) + .isInstanceOf(AccessDeniedException.class); + } + + @Test + @DisplayName("CONNECT가 아닌 프레임은 검증 없이 통과시킨다") + void preSend_skipsValidation_forNonConnectFrames() { + // Given + StompHeaderAccessor accessor = StompHeaderAccessor.create(StompCommand.SEND); + accessor.setLeaveMutable(true); + Message sendMessage = MessageBuilder.createMessage(new byte[0], accessor.getMessageHeaders()); + + // When + Message result = interceptor.preSend(sendMessage, channel); + + // Then + assertThat(result).isSameAs(sendMessage); + then(jwtProvider).shouldHaveNoInteractions(); + } + + private Message connectMessage(String authorizationHeader) { + StompHeaderAccessor accessor = StompHeaderAccessor.create(StompCommand.CONNECT); + if (authorizationHeader != null) { + accessor.setNativeHeader("Authorization", authorizationHeader); + } + accessor.setLeaveMutable(true); + return MessageBuilder.createMessage(new byte[0], accessor.getMessageHeaders()); + } +} diff --git a/src/test/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardWebSocketIntegrationTest.java b/src/test/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardWebSocketIntegrationTest.java new file mode 100644 index 0000000..5801eee --- /dev/null +++ b/src/test/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardWebSocketIntegrationTest.java @@ -0,0 +1,210 @@ +package com.opensource.docgrid.domain.dashboard.websocket; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.lang.reflect.Type; +import java.util.List; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.web.server.LocalServerPort; +import org.springframework.messaging.converter.MappingJackson2MessageConverter; +import org.springframework.messaging.simp.stomp.StompCommand; +import org.springframework.messaging.simp.stomp.StompFrameHandler; +import org.springframework.messaging.simp.stomp.StompHeaders; +import org.springframework.messaging.simp.stomp.StompSession; +import org.springframework.messaging.simp.stomp.StompSessionHandlerAdapter; +import org.springframework.test.context.ActiveProfiles; +import org.springframework.test.context.DynamicPropertyRegistry; +import org.springframework.test.context.DynamicPropertySource; +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.dashboard.controller.DashboardWebSocketController; +import com.opensource.docgrid.domain.dashboard.dto.response.DashboardSummaryResponse; +import com.opensource.docgrid.domain.dashboard.dto.response.DocumentsSummaryResponse; +import com.opensource.docgrid.domain.dashboard.dto.response.JobsSummaryResponse; +import com.opensource.docgrid.domain.dashboard.dto.response.SearchSummaryResponse; +import com.opensource.docgrid.domain.dashboard.dto.response.WorkersSummaryResponse; + +/** + * 실제 STOMP Client로 {@code /ws} 연결부터 {@code /topic/dashboard} 구독, 수동 push 수신까지 + * 관통하는 통합 테스트. + * + *

CONNECT 시점 JWT 검증({@code StompAuthChannelInterceptor})과 SUBSCRIBE 시점 ADMIN 권한 + * 검증({@code DashboardSubscriptionAuthorizationInterceptor})이 실제 Channel Interceptor + * 체인에서 함께 동작하는지 확인한다. + */ +@Tag("integration") +@ActiveProfiles("test") +@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT) +@DisplayName("Dashboard WebSocket 실시간 push 통합 테스트") +class DashboardWebSocketIntegrationTest { + + private static final long TIMEOUT_SECONDS = 5; + + @LocalServerPort + private int port; + + @Autowired + private JwtProvider jwtProvider; + + @Autowired + private DashboardWebSocketController dashboardWebSocketController; + + private WebSocketStompClient stompClient; + + @DynamicPropertySource + static void configureJwt(DynamicPropertyRegistry registry) { + registry.add("jwt.secret", () -> "docgrid-dashboard-websocket-integration-test-secret-key-2026"); + } + + @BeforeEach + void setUp() { + stompClient = new WebSocketStompClient(new StandardWebSocketClient()); + stompClient.setMessageConverter(new MappingJackson2MessageConverter()); + } + + @Test + @DisplayName("정상 케이스: ADMIN 토큰으로 구독하면 sendDashboardUpdate() push를 즉시 수신한다") + void receivesPush_whenAdminSubscribed() throws Exception { + // Given + BlockingQueue failures = new LinkedBlockingQueue<>(); + StompSession session = connect(adminToken(), failures); + BlockingQueue received = new LinkedBlockingQueue<>(); + subscribeDashboard(session, received); + + // When + DashboardSummaryResponse summary = sampleSummary(); + dashboardWebSocketController.sendDashboardUpdate(summary); + + // Then + DashboardSummaryResponse result = received.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS); + assertThat(result).isNotNull(); + assertThat(result.documents().total()).isEqualTo(summary.documents().total()); + assertThat(result.jobs().failed()).isEqualTo(summary.jobs().failed()); + assertThat(failures).isEmpty(); + + session.disconnect(); + } + + @Test + @DisplayName("예외 케이스: 토큰 없이 CONNECT하면 연결이 거부된다") + void rejectsConnect_whenNoToken() throws Exception { + // Given + BlockingQueue failures = new LinkedBlockingQueue<>(); + StompHeaders connectHeaders = new StompHeaders(); + + // When + StompSession session = null; + Throwable connectFailure = null; + try { + session = stompClient + .connectAsync(wsUrl(), (WebSocketHttpHeaders) null, connectHeaders, failureCapturingHandler(failures)) + .get(TIMEOUT_SECONDS, TimeUnit.SECONDS); + } catch (ExecutionException e) { + connectFailure = e.getCause(); + } + + // Then — Transport 레벨에서 바로 실패하거나, 연결은 되고 STOMP 레벨 ERROR로 이어진다. + if (session == null) { + assertThat(connectFailure).isNotNull(); + } else { + assertThat(failures.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isNotNull(); + } + } + + @Test + @DisplayName("예외 케이스: ADMIN이 아닌 사용자는 /topic/dashboard SUBSCRIBE가 거부된다") + void rejectsSubscribe_whenNotAdmin() throws Exception { + // Given + BlockingQueue failures = new LinkedBlockingQueue<>(); + StompSession session = connect(userToken(), failures); + + // When + session.subscribe("/topic/dashboard", new StompFrameHandler() { + @Override + public Type getPayloadType(StompHeaders headers) { + return DashboardSummaryResponse.class; + } + + @Override + public void handleFrame(StompHeaders headers, Object payload) { + // 정상 케이스가 아니므로 페이로드를 받으면 안 된다. + } + }); + + // Then — SUBSCRIBE 거부 시 서버가 세션을 이미 닫으므로 별도 disconnect()는 호출하지 않는다. + assertThat(failures.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isNotNull(); + } + + private StompSession connect(String token, BlockingQueue failures) throws Exception { + StompHeaders connectHeaders = new StompHeaders(); + connectHeaders.add("Authorization", "Bearer " + token); + + return stompClient + .connectAsync(wsUrl(), (WebSocketHttpHeaders) null, connectHeaders, failureCapturingHandler(failures)) + .get(TIMEOUT_SECONDS, TimeUnit.SECONDS); + } + + private StompSessionHandlerAdapter failureCapturingHandler(BlockingQueue failures) { + return new StompSessionHandlerAdapter() { + @Override + public void handleException( + StompSession session, StompCommand command, StompHeaders headers, byte[] payload, Throwable exception + ) { + failures.add(exception); + } + + @Override + public void handleTransportError(StompSession session, Throwable exception) { + failures.add(exception); + } + }; + } + + private void subscribeDashboard(StompSession session, BlockingQueue received) { + session.subscribe("/topic/dashboard", new StompFrameHandler() { + @Override + public Type getPayloadType(StompHeaders headers) { + return DashboardSummaryResponse.class; + } + + @Override + public void handleFrame(StompHeaders headers, Object payload) { + received.add((DashboardSummaryResponse) payload); + } + }); + } + + private String wsUrl() { + return "ws://localhost:" + port + "/ws/websocket"; + } + + private String adminToken() { + return jwtProvider.generateToken(1L, "dashboard-admin@example.com", List.of("ADMIN")); + } + + private String userToken() { + return jwtProvider.generateToken(2L, "dashboard-user@example.com", List.of("USER")); + } + + private DashboardSummaryResponse sampleSummary() { + return new DashboardSummaryResponse( + new DocumentsSummaryResponse(25368L, 21742L, 132L), + new JobsSummaryResponse(132L, 8L, 27L, 3200L), + new WorkersSummaryResponse(5L, 6L), + new SearchSummaryResponse(342L) + ); + } +} From f716fc9b7d0b08c49e843d81bdd2e4472e1c5d69 Mon Sep 17 00:00:00 2001 From: kangcheolung Date: Mon, 10 Aug 2026 19:12:10 +0900 Subject: [PATCH 9/9] =?UTF-8?q?fix:=20#137=20=EC=BD=94=EB=93=9C=EB=9E=98?= =?UTF-8?q?=EB=B9=97=20=EB=A6=AC=EB=B7=B0=20=EB=B0=98=EC=98=81=20-=20/topi?= =?UTF-8?q?c/dashboard=20SEND=20=EC=B0=A8=EB=8B=A8,=20=EC=88=9C=EC=84=9C?= =?UTF-8?q?=20=EC=A3=BC=EC=84=9D=20=EB=B3=B4=EA=B0=95,=20=ED=85=8C?= =?UTF-8?q?=EC=8A=A4=ED=8A=B8=20=EB=8F=99=EA=B8=B0=ED=99=94=20=EB=B0=A9?= =?UTF-8?q?=EC=8B=9D=20=EA=B5=90=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Sonnet 5 --- ...ng-#137-ragops-dashboard-websocket-push.md | 64 ++++++++++++----- ...dSubscriptionAuthorizationInterceptor.java | 22 ++++-- .../docgrid/global/config/SecurityConfig.java | 8 ++- .../global/config/WebSocketConfig.java | 4 +- .../DashboardWebSocketIntegrationTest.java | 71 ++++++++++++++++--- 5 files changed, 134 insertions(+), 35 deletions(-) diff --git a/docs/design/kangcheolung-#137-ragops-dashboard-websocket-push.md b/docs/design/kangcheolung-#137-ragops-dashboard-websocket-push.md index 0bce997..3d22a60 100644 --- a/docs/design/kangcheolung-#137-ragops-dashboard-websocket-push.md +++ b/docs/design/kangcheolung-#137-ragops-dashboard-websocket-push.md @@ -16,7 +16,7 @@ push를 트리거할지(재처리 클릭, Worker 상태 전이)는 후속 이슈 - `/ws`(SockJS) 연결과 `/topic/dashboard` 구독이 정상 동작한다. - 상태 변경 시점에만 push하며, 고정 주기 폴링을 쓰지 않는다. -- ADMIN이 아닌 사용자는 연결 또는 구독 시점에 거부된다. +- ADMIN이 아닌 사용자는 `/topic/dashboard` 구독(SUBSCRIBE) 시 거부된다 (CONNECT는 유효한 JWT만 있으면 통과한다). - 이 앱의 JWT 인증(stateless, 세션 없음) 모델과 충돌 없이 동작한다. ## 2. 범위 @@ -156,6 +156,9 @@ public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { @Override public void configureClientInboundChannel(ChannelRegistration registration) { + // 1. StompAuthChannelInterceptor가 CONNECT 프레임의 JWT를 검증하고 세션에 Principal을 부착한다. + // 2. DashboardSubscriptionAuthorizationInterceptor가 그 Principal로 SUBSCRIBE·SEND 권한을 검증한다. + // 순서가 바뀌면 2번 시점에 Principal이 아직 없어 항상 거부된다. registration.interceptors(stompAuthChannelInterceptor, dashboardSubscriptionAuthorizationInterceptor); } } @@ -245,7 +248,7 @@ WebSocket 자체는 그냥 양방향 파이프고 "이게 연결 요청인지 `SecurityContextHolder`를 새로 채우지만, WebSocket은 연결이 오래 유지되는 세션이라 이렇게 세션 레벨에 신원을 붙여두고 이후 모든 프레임에서 재사용한다. -### 5.5 `DashboardSubscriptionAuthorizationInterceptor.java` — SUBSCRIBE 시점 ROLE_ADMIN 인가 +### 5.5 `DashboardSubscriptionAuthorizationInterceptor.java` — `/topic/dashboard` SUBSCRIBE·SEND 인가 ```java @Component @@ -258,10 +261,15 @@ public class DashboardSubscriptionAuthorizationInterceptor implements ChannelInt public Message preSend(Message message, MessageChannel channel) { StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); - if (accessor != null - && StompCommand.SUBSCRIBE.equals(accessor.getCommand()) - && DASHBOARD_TOPIC.equals(accessor.getDestination()) - && !isAdmin(accessor.getUser())) { + if (accessor == null || !DASHBOARD_TOPIC.equals(accessor.getDestination())) { + return message; + } + + if (StompCommand.SEND.equals(accessor.getCommand())) { + throw new AccessDeniedException("이 목적지로는 메시지를 보낼 수 없습니다."); + } + + if (StompCommand.SUBSCRIBE.equals(accessor.getCommand()) && !isAdmin(accessor.getUser())) { throw new AccessDeniedException("대시보드 구독 권한이 없습니다."); } @@ -279,13 +287,19 @@ public class DashboardSubscriptionAuthorizationInterceptor implements ChannelInt } ``` -5.4가 "누구냐"를 확인했다면, 이 클래스는 "이 사람이 대시보드를 볼 자격이 있냐"만 본다. 5.4를 -통과해서 이미 세션이 열려있는(=로그인 확인된) 사람 중에서, `/topic/dashboard`를 구독하려 할 -때 ADMIN인지만 한 번 더 확인한다 — CONNECT 검증 하나에만 의존하지 않는 이중 방어다. +5.4가 "누구냐"를 확인했다면, 이 클래스는 "이 사람이 `/topic/dashboard`에 뭘 할 자격이 있냐"를 +본다. 목적지가 `/topic/dashboard`가 아니면 바로 통과시키고, 그 목적지에 한해서만 두 커맨드를 +따로 판단한다: -`preSend`의 `if` 조건 4개가 전부 참이어야 거부한다: ① Accessor가 null이 아니고 ② SUBSCRIBE -프레임이고 ③ 목적지가 `/topic/dashboard`고 ④ `isAdmin()`이 `false`일 때. 넷 중 하나라도 아니면 -그냥 통과시킨다. +- **SEND 차단(무조건)** — `enableSimpleBroker("/topic")` 구성에서는 클라이언트가 `/topic/dashboard`로 + STOMP SEND 프레임을 보내면 SimpleBroker가 그걸 그대로 구독자 전원에게 브로드캐스트한다. 즉 + CONNECT만 통과한 일반 사용자도 서버인 척 위조 지표를 ADMIN 구독자에게 보낼 수 있다는 뜻이다. + 실제 push는 `DashboardWebSocketController`가 `clientInboundChannel`을 거치지 않는 + `SimpMessagingTemplate`으로만 하므로, 클라이언트발 SEND는 ADMIN 여부와 무관하게 전부 차단해도 + 정상 기능에 영향이 없다. 초기 구현에는 이 분기가 없어서 CodeRabbit 리뷰로 뒤늦게 발견했다 + (7장 참고). +- **SUBSCRIBE는 ADMIN만 허용** — 5.4를 통과해서 이미 세션이 열려있는(=로그인 확인된) 사람 중에서, + ADMIN인지만 한 번 더 확인한다. CONNECT 검증 하나에만 의존하지 않는 이중 방어다. `isAdmin`은 `accessor.getUser()`(타입은 `java.security.Principal`)를 `Authentication`으로 다운캐스트해서 `getAuthorities()`를 봐야 권한 목록에 접근할 수 있다(`Principal` 자체는 이름만 @@ -300,9 +314,11 @@ public class DashboardSubscriptionAuthorizationInterceptor implements ChannelInt ```diff .requestMatchers("/auth/signup", "/auth/login").permitAll() -+ // WebSocket 핸드셰이크는 여기서 인증하지 않는다. 네이티브 websocket Transport는 -+ // Upgrade 요청에 커스텀 헤더를 실을 수 없어, 인증은 StompAuthChannelInterceptor가 -+ // STOMP CONNECT 프레임에서 담당하고 목적지별 인가는 DashboardSubscriptionAuthorizationInterceptor가 담당한다. ++ // WebSocket 인증·인가는 3단계로 나뉜다. 네이티브 websocket Transport가 Upgrade ++ // 요청에 커스텀 헤더를 못 실어서, 여기(HTTP)에서는 검증하지 않는다: ++ // 1. HTTP 핸드셰이크(여기) — permitAll ++ // 2. STOMP CONNECT — StompAuthChannelInterceptor가 JWT 검증 ++ // 3. STOMP SUBSCRIBE·SEND — DashboardSubscriptionAuthorizationInterceptor가 목적지별 권한 검증 + .requestMatchers("/ws/**").permitAll() .requestMatchers("/admin/**").hasRole("ADMIN") .anyRequest().authenticated() @@ -365,9 +381,21 @@ Worker 상태 전이 이벤트 훅)에서 채운다. `@SpringBootTest(webEnvironment = RANDOM_PORT)` + 실제 `WebSocketStompClient`로 서버에 직접 연결해서 검증한다. -- ADMIN 토큰으로 구독 → `sendDashboardUpdate()` 호출 시 1초 내 수신 +- ADMIN 토큰으로 구독 → `sendDashboardUpdate()` 호출 시 1초 내 수신 (완료 기준의 delivery + SLA를 그대로 타임아웃 값으로 사용) - 토큰 없이 CONNECT → 연결 거부 (Transport 레벨 실패 또는 STOMP ERROR 프레임 중 어느 쪽이든 대응) - USER(비 ADMIN) 토큰 → CONNECT는 성공하지만 `/topic/dashboard` SUBSCRIBE는 거부 +- USER 토큰으로 `/topic/dashboard`에 직접 SEND → 거부되고, 동시에 구독 중인 ADMIN에게도 + 전달되지 않음 (5.5절 SEND 차단 검증) + +**구독 완료 대기는 STOMP Receipt가 아니라 `SimpUserRegistry` 폴링으로 한다.** 원래는 +`StompSession.Subscription.addReceiptTask()`로 브로커가 SUBSCRIBE를 처리했다는 RECEIPT +프레임을 기다리려 했는데, `enableSimpleBroker`(in-memory `SimpleBrokerMessageHandler`)는 +STOMP Receipt를 아예 구현하지 않는다 — `DefaultStompSession$ReceiptHandler`는 클라이언트 쪽 +로직일 뿐이라 서버가 RECEIPT를 보내는 기능 자체가 없어 영원히 대기하다 타임아웃났다. 대신 +같은 JVM에서 실행 중인 `SimpUserRegistry.findSubscriptions(...)`로 서버가 실제로 그 구독을 +인지했는지 직접 확인하고, 확인될 때까지 짧은 간격(20ms)으로 재폴링한다. 고정 sleep 하나로 +"아마 됐겠지"하고 넘어가지 않고, 서버의 실제 상태를 근거로 대기한다. ## 7. 오류 계약 @@ -376,6 +404,7 @@ Worker 상태 전이 이벤트 훅)에서 채운다. | `Authorization` 헤더 없이 CONNECT | `StompAuthChannelInterceptor`가 `AccessDeniedException` → 연결 종료 | | 토큰 만료·서명 무효 | 위와 동일 | | ADMIN이 아닌 사용자가 `/topic/dashboard` SUBSCRIBE | `DashboardSubscriptionAuthorizationInterceptor`가 `AccessDeniedException` → 구독 거부, 세션 종료 | +| 아무 클라이언트가 `/topic/dashboard`로 SEND | `DashboardSubscriptionAuthorizationInterceptor`가 `AccessDeniedException` → SEND 거부, ADMIN 구독자에게도 전달 안 됨 (CodeRabbit 리뷰로 발견 — 초기 구현엔 이 분기가 없어 인증만 된 일반 사용자가 위조 지표를 ADMIN에게 보낼 수 있었음) | | `/ws/**` HTTP 핸드셰이크 자체 | 항상 permitAll, 여기서는 거부되지 않음 | ## 8. 커밋 분할 @@ -392,4 +421,5 @@ Worker 상태 전이 이벤트 훅)에서 채운다. - `sendDashboardUpdate()` 호출 시 구독 중인 클라이언트가 1초 내 최신 지표를 수신한다 - 토큰 없음/무효 토큰으로 CONNECT 시 거부된다 - ADMIN이 아닌 사용자는 `/topic/dashboard` SUBSCRIBE가 거부된다 -- 전체 빌드(`./gradlew build`)가 회귀 없이 통과한다 (733개 테스트, failures 0, errors 0) +- 클라이언트가 `/topic/dashboard`로 SEND하면 거부되고, 구독 중인 ADMIN에게도 전달되지 않는다 +- 전체 빌드(`./gradlew build`)가 회귀 없이 통과한다 (734개 테스트, failures 0, errors 0) diff --git a/src/main/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardSubscriptionAuthorizationInterceptor.java b/src/main/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardSubscriptionAuthorizationInterceptor.java index 3930e10..11390eb 100644 --- a/src/main/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardSubscriptionAuthorizationInterceptor.java +++ b/src/main/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardSubscriptionAuthorizationInterceptor.java @@ -14,11 +14,18 @@ import org.springframework.stereotype.Component; /** - * {@code /topic/dashboard} SUBSCRIBE 요청에 ROLE_ADMIN 권한을 요구한다. + * {@code /topic/dashboard} 목적지를 ADMIN 전용으로 보호한다. * *

{@code StompAuthChannelInterceptor}가 CONNECT 시점에 세션에 부착한 Principal을 재사용해 * 목적지 접근 시점에 다시 한 번 검증한다. CONNECT 검증 하나에만 의존하지 않는 이중 방어다. * + *

SUBSCRIBE뿐 아니라 SEND도 차단한다. {@code enableSimpleBroker("/topic")} 구성에서는 + * 클라이언트가 {@code /topic/dashboard}로 STOMP SEND 프레임을 보내면 SimpleBroker가 이를 그대로 + * 구독자 전원에게 브로드캐스트한다 — 인증만 된 일반 사용자도 위조된 지표를 ADMIN 구독자에게 보낼 + * 수 있다는 뜻이다. 실제 push는 {@code DashboardWebSocketController}가 {@code clientInboundChannel}을 + * 거치지 않는 {@code SimpMessagingTemplate}으로만 하므로, 이 목적지로의 클라이언트발 SEND는 + * ADMIN 여부와 무관하게 전부 차단해도 정상 기능에 영향이 없다. + * *

{@code @EnableWebSocketSecurity}(Spring Security 메시지 인가 DSL)는 STOMP endpoint가 * 등록된 것을 감지하면 세션 기반 CSRF 토큰을 무조건 요구하는 {@code CsrfChannelInterceptor}를 * 함께 붙인다. 이 앱은 세션이 없는 stateless JWT 인증이라 CSRF 토큰이 존재할 수 없어 모든 CONNECT가 @@ -34,10 +41,15 @@ public class DashboardSubscriptionAuthorizationInterceptor implements ChannelInt public Message preSend(Message message, MessageChannel channel) { StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); - if (accessor != null - && StompCommand.SUBSCRIBE.equals(accessor.getCommand()) - && DASHBOARD_TOPIC.equals(accessor.getDestination()) - && !isAdmin(accessor.getUser())) { + if (accessor == null || !DASHBOARD_TOPIC.equals(accessor.getDestination())) { + return message; + } + + if (StompCommand.SEND.equals(accessor.getCommand())) { + throw new AccessDeniedException("이 목적지로는 메시지를 보낼 수 없습니다."); + } + + if (StompCommand.SUBSCRIBE.equals(accessor.getCommand()) && !isAdmin(accessor.getUser())) { throw new AccessDeniedException("대시보드 구독 권한이 없습니다."); } diff --git a/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java b/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java index 5c5032a..c2e522c 100644 --- a/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java +++ b/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java @@ -39,9 +39,11 @@ public SecurityFilterChain filterChain(HttpSecurity http) throws Exception { .requestMatchers("/swagger-ui/**", "/v3/api-docs/**").permitAll() .requestMatchers("/departments").permitAll() .requestMatchers("/auth/signup", "/auth/login").permitAll() - // WebSocket 핸드셰이크는 여기서 인증하지 않는다. 네이티브 websocket Transport는 - // Upgrade 요청에 커스텀 헤더를 실을 수 없어, 인증은 StompAuthChannelInterceptor가 - // STOMP CONNECT 프레임에서 담당하고 목적지별 인가는 DashboardSubscriptionAuthorizationInterceptor가 담당한다. + // WebSocket 인증·인가는 3단계로 나뉜다. 네이티브 websocket Transport가 Upgrade + // 요청에 커스텀 헤더를 못 실어서, 여기(HTTP)에서는 검증하지 않는다: + // 1. HTTP 핸드셰이크(여기) — permitAll + // 2. STOMP CONNECT — StompAuthChannelInterceptor가 JWT 검증 + // 3. STOMP SUBSCRIBE·SEND — DashboardSubscriptionAuthorizationInterceptor가 목적지별 권한 검증 .requestMatchers("/ws/**").permitAll() .requestMatchers("/admin/**").hasRole("ADMIN") .anyRequest().authenticated() diff --git a/src/main/java/com/opensource/docgrid/global/config/WebSocketConfig.java b/src/main/java/com/opensource/docgrid/global/config/WebSocketConfig.java index 619dd04..34bb1e6 100644 --- a/src/main/java/com/opensource/docgrid/global/config/WebSocketConfig.java +++ b/src/main/java/com/opensource/docgrid/global/config/WebSocketConfig.java @@ -41,7 +41,9 @@ public void configureMessageBroker(MessageBrokerRegistry registry) { @Override public void configureClientInboundChannel(ChannelRegistration registration) { - // CONNECT 인증이 SUBSCRIBE 인가보다 먼저 Principal을 세션에 부착해야 하므로 순서를 고정한다. + // 1. StompAuthChannelInterceptor가 CONNECT 프레임의 JWT를 검증하고 세션에 Principal을 부착한다. + // 2. DashboardSubscriptionAuthorizationInterceptor가 그 Principal로 SUBSCRIBE·SEND 권한을 검증한다. + // 순서가 바뀌면 2번 시점에 Principal이 아직 없어 항상 거부된다. registration.interceptors(stompAuthChannelInterceptor, dashboardSubscriptionAuthorizationInterceptor); } } diff --git a/src/test/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardWebSocketIntegrationTest.java b/src/test/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardWebSocketIntegrationTest.java index 5801eee..48bc081 100644 --- a/src/test/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardWebSocketIntegrationTest.java +++ b/src/test/java/com/opensource/docgrid/domain/dashboard/websocket/DashboardWebSocketIntegrationTest.java @@ -22,6 +22,7 @@ import org.springframework.messaging.simp.stomp.StompHeaders; import org.springframework.messaging.simp.stomp.StompSession; import org.springframework.messaging.simp.stomp.StompSessionHandlerAdapter; +import org.springframework.messaging.simp.user.SimpUserRegistry; import org.springframework.test.context.ActiveProfiles; import org.springframework.test.context.DynamicPropertyRegistry; import org.springframework.test.context.DynamicPropertySource; @@ -41,8 +42,8 @@ * 실제 STOMP Client로 {@code /ws} 연결부터 {@code /topic/dashboard} 구독, 수동 push 수신까지 * 관통하는 통합 테스트. * - *

CONNECT 시점 JWT 검증({@code StompAuthChannelInterceptor})과 SUBSCRIBE 시점 ADMIN 권한 - * 검증({@code DashboardSubscriptionAuthorizationInterceptor})이 실제 Channel Interceptor + *

CONNECT 시점 JWT 검증({@code StompAuthChannelInterceptor})과 SUBSCRIBE·SEND 시점 ADMIN + * 권한 검증({@code DashboardSubscriptionAuthorizationInterceptor})이 실제 Channel Interceptor * 체인에서 함께 동작하는지 확인한다. */ @Tag("integration") @@ -51,7 +52,10 @@ @DisplayName("Dashboard WebSocket 실시간 push 통합 테스트") class DashboardWebSocketIntegrationTest { + private static final String DASHBOARD_TOPIC = "/topic/dashboard"; private static final long TIMEOUT_SECONDS = 5; + private static final long DELIVERY_TIMEOUT_SECONDS = 1; + private static final long SUBSCRIPTION_POLL_INTERVAL_MILLIS = 20; @LocalServerPort private int port; @@ -62,6 +66,9 @@ class DashboardWebSocketIntegrationTest { @Autowired private DashboardWebSocketController dashboardWebSocketController; + @Autowired + private SimpUserRegistry simpUserRegistry; + private WebSocketStompClient stompClient; @DynamicPropertySource @@ -76,20 +83,20 @@ void setUp() { } @Test - @DisplayName("정상 케이스: ADMIN 토큰으로 구독하면 sendDashboardUpdate() push를 즉시 수신한다") + @DisplayName("정상 케이스: ADMIN 토큰으로 구독하면 sendDashboardUpdate() push를 1초 이내 수신한다") void receivesPush_whenAdminSubscribed() throws Exception { // Given BlockingQueue failures = new LinkedBlockingQueue<>(); StompSession session = connect(adminToken(), failures); BlockingQueue received = new LinkedBlockingQueue<>(); - subscribeDashboard(session, received); + subscribeDashboardAndAwaitRegistration(session, received); // When DashboardSummaryResponse summary = sampleSummary(); dashboardWebSocketController.sendDashboardUpdate(summary); - // Then - DashboardSummaryResponse result = received.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS); + // Then — 완료 기준(1초 이내 전달)을 그대로 타임아웃으로 사용해 SLA를 검증한다. + DashboardSummaryResponse result = received.poll(DELIVERY_TIMEOUT_SECONDS, TimeUnit.SECONDS); assertThat(result).isNotNull(); assertThat(result.documents().total()).isEqualTo(summary.documents().total()); assertThat(result.jobs().failed()).isEqualTo(summary.jobs().failed()); @@ -132,7 +139,7 @@ void rejectsSubscribe_whenNotAdmin() throws Exception { StompSession session = connect(userToken(), failures); // When - session.subscribe("/topic/dashboard", new StompFrameHandler() { + session.subscribe(DASHBOARD_TOPIC, new StompFrameHandler() { @Override public Type getPayloadType(StompHeaders headers) { return DashboardSummaryResponse.class; @@ -148,6 +155,27 @@ public void handleFrame(StompHeaders headers, Object payload) { assertThat(failures.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isNotNull(); } + @Test + @DisplayName("예외 케이스: /topic/dashboard로 SEND하면 거부되고 구독자에게 전달되지 않는다") + void rejectsSend_toDashboardTopic() throws Exception { + // Given — ADMIN이 정상 구독 중인 상태를 먼저 만든다. + BlockingQueue adminFailures = new LinkedBlockingQueue<>(); + StompSession adminSession = connect(adminToken(), adminFailures); + BlockingQueue received = new LinkedBlockingQueue<>(); + subscribeDashboardAndAwaitRegistration(adminSession, received); + + // When — 별도 세션(ADMIN 아님)이 서버인 척 위조 페이로드를 직접 SEND한다. + BlockingQueue senderFailures = new LinkedBlockingQueue<>(); + StompSession senderSession = connect(userToken(), senderFailures); + senderSession.send(DASHBOARD_TOPIC, sampleSummary()); + + // Then — SEND 자체가 거부되고, SimpleBroker가 구독자에게 브로드캐스트하지 않는다. + assertThat(senderFailures.poll(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isNotNull(); + assertThat(received.poll(DELIVERY_TIMEOUT_SECONDS, TimeUnit.SECONDS)).isNull(); + + adminSession.disconnect(); + } + private StompSession connect(String token, BlockingQueue failures) throws Exception { StompHeaders connectHeaders = new StompHeaders(); connectHeaders.add("Authorization", "Bearer " + token); @@ -173,8 +201,21 @@ public void handleTransportError(StompSession session, Throwable exception) { }; } - private void subscribeDashboard(StompSession session, BlockingQueue received) { - session.subscribe("/topic/dashboard", new StompFrameHandler() { + /** + * 구독을 요청하고, 서버 브로커의 실제 구독 registry에 등록될 때까지 폴링으로 대기한다. + * + *

{@code enableSimpleBroker}(in-memory {@code SimpleBrokerMessageHandler})는 STOMP + * Receipt를 구현하지 않는다 — {@code StompSession.Subscription.addReceiptTask()}로는 서버가 + * 절대 RECEIPT 프레임을 보내주지 않아 영원히 대기하게 된다. 대신 같은 JVM에서 실행 중인 + * {@link SimpUserRegistry}로 서버가 실제로 이 구독을 인지했는지 직접 확인한다. 구독 직후 + * 곧바로 push하면 브로커가 SUBSCRIBE 등록을 마치기 전에 push가 먼저 도착해 유실될 수 있어, + * 고정 sleep 대신 실제 서버 상태가 확정될 때까지 짧은 간격으로 재확인한다. + */ + private void subscribeDashboardAndAwaitRegistration( + StompSession session, + BlockingQueue received + ) throws InterruptedException { + session.subscribe(DASHBOARD_TOPIC, new StompFrameHandler() { @Override public Type getPayloadType(StompHeaders headers) { return DashboardSummaryResponse.class; @@ -185,6 +226,18 @@ public void handleFrame(StompHeaders headers, Object payload) { received.add((DashboardSummaryResponse) payload); } }); + + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(TIMEOUT_SECONDS); + while (System.nanoTime() < deadline) { + boolean registered = !simpUserRegistry + .findSubscriptions(subscription -> DASHBOARD_TOPIC.equals(subscription.getDestination())) + .isEmpty(); + if (registered) { + return; + } + Thread.sleep(SUBSCRIPTION_POLL_INTERVAL_MILLIS); + } + throw new AssertionError("구독이 " + TIMEOUT_SECONDS + "초 안에 서버에 등록되지 않았습니다."); } private String wsUrl() {