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
@@ -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;
Expand All @@ -12,7 +14,7 @@
import lombok.RequiredArgsConstructor;

/**
* HTTP 관리자 인가에 필요한 역할을 Redis나 standby 없이 현재 primary에서 확인한다.
* HTTP 관리자 인가와 대시보드 push에 필요한 역할을 Redis나 standby 없이 현재 primary에서 확인한다.
*
* <p>호출자의 read-only 트랜잭션을 상속하지 않으며, OpenProxy가 잘못 라우팅하면
* 권한을 추정하지 않고 요청을 실패시킨다. 일반 API·WebSocket의 역할 캐시는 담당하지 않는다.
Expand All @@ -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<String> findCurrentRoles(Long userId) {
// 1. 명시적 read-write 트랜잭션 안에서 OpenProxy의 실제 도착 역할을 확인한다.
requirePrimary();

// 2. 같은 트랜잭션의 최신 primary에서 역할을 조회한다.
return userRoleRepository.findRoleCodesByUserId(userId);
}

@Transactional(propagation = Propagation.REQUIRES_NEW)
public Set<Long> findCurrentAdminUserIds(List<Long> userIds) {
if (userIds.isEmpty()) {
return Set.of();
}

// 1. 단일 read-write 트랜잭션에서 primary를 확인해 standby의 낡은 ADMIN을 쓰지 않는다.
requirePrimary();

// 2. 큰 세션 집합도 크기가 제한된 IN 쿼리로 읽되, 결과는 이 push에만 쓰는 불변 집합으로 만든다.
Set<Long> 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);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
package com.opensource.docgrid.domain.auth.websocket;

import java.util.Map;
import java.util.Set;

/**
* 대시보드 push 한 번에 대해 후보 물리 세션과 primary에서 확인한 ADMIN 사용자만 내부 메시지에 전달한다.
*
* <p>브로커가 복제한 수신자별 MESSAGE에서만 읽으며 클라이언트 native STOMP 헤더에는 싣지 않는다.
* 다른 push에 재사용하지 않고, 로그에 사용자 ID가 노출되지 않도록 문자열 표현도 제한한다.
*/
public record DashboardAuthorizationSnapshot(Map<String, Long> candidateSessions, Set<Long> 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]";
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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 판정으로 제한한다.
*
* <p>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) {
Expand All @@ -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;
}
}
Original file line number Diff line number Diff line change
@@ -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;
Expand All @@ -12,7 +22,9 @@
*
* <p>이 컨트롤러는 직접 지표를 집계하지 않는다. 호출하는 쪽(재처리 트리거, 상태 전이 이벤트
* 리스너 등)이 {@code DashboardQueryService}로 최신 Snapshot을 계산해 넘겨주면 그대로
* 브로드캐스트만 한다.
* 브로드캐스트한다. 호출자가 요약을 먼저 계산한 뒤 이 메서드에서 권한을 판정해야,
* 회수 전 판정으로 회수 후 새로 계산한 요약을 승인하지 않는다. 판정은 push 사이에 재사용하지 않는다.
* 구독자별 DB 조회를 반복하지 않도록 각 push의 ADMIN 후보를 primary에서 일괄 확인한다.
*/
@Component
@RequiredArgsConstructor
Expand All @@ -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<String, Long> candidateSessions = stompSessionRegistry.authenticatedSessions().stream()
.filter(session -> session.authorization().roles().contains("ADMIN"))
.collect(Collectors.toUnmodifiableMap(session -> session.sessionId(),
session -> session.authorization().userId()));
List<Long> 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());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,8 @@

/**
* 관리자 재처리 버튼 클릭을 A 담당자의 {@link EmbeddingJobManualRetryService}로 위임하고,
* 성공 시 최신 대시보드 집계를 WebSocket으로 push하는 Command Service.
* 성공 시 최신 대시보드 집계를 WebSocket으로 push하는 Command Service. 알림 인가 조회가
* 실패해도 이미 커밋된 재처리 결과를 HTTP 실패로 잘못 보고하지 않는다.
*
* <p>{@code embedding_jobs}를 직접 update하지 않는다 — 상태 전환은 전부 A의 Service를 경유한다.
* FAILED 목록 조회만 {@link EmbeddingJobRepository}를 직접 읽는다(쓰기가 아니므로 A/B 경계 위반이
Expand Down Expand Up @@ -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;
}

Expand Down Expand Up @@ -82,7 +83,7 @@ public RetryAllJobsResponse retryAllFailedJobs() {

// 3. 실제로 바뀐 게 있을 때만(1건 이상 성공) push한다 — 전부 실패하면 push할 변경사항이 없다.
if (retriedCount > 0) {
dashboardWebSocketController.sendDashboardUpdate(dashboardQueryService.getSummary());
pushDashboardUpdateAfterRetry();
}

return new RetryAllJobsResponse(
Expand All @@ -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());
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,10 @@ public interface UserRoleRepository extends JpaRepository<UserRole, Long> {
@Query("SELECT ur.role.code FROM UserRole ur WHERE ur.user.id = :userId")
List<String> 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<Long> findAdminUserIdsByUserIdIn(@Param("userIds") List<Long> userIds);

boolean existsByUserIdAndRoleCode(Long userId, String roleCode);

Optional<UserRole> findByUserIdAndRoleCode(Long userId, String roleCode);
Expand Down
Loading
Loading