diff --git a/backend/src/main/java/com/opensource/docgrid/domain/auth/jwt/JwtAuthenticationFilter.java b/backend/src/main/java/com/opensource/docgrid/domain/auth/jwt/JwtAuthenticationFilter.java
index cbe14f94..2d6761f2 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/auth/jwt/JwtAuthenticationFilter.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/auth/jwt/JwtAuthenticationFilter.java
@@ -3,12 +3,18 @@
import java.io.IOException;
import java.util.List;
+import org.springframework.http.MediaType;
import org.springframework.security.authentication.UsernamePasswordAuthenticationToken;
import org.springframework.security.core.authority.SimpleGrantedAuthority;
import org.springframework.security.core.context.SecurityContextHolder;
+import org.springframework.security.web.util.matcher.RequestMatcher;
import org.springframework.util.StringUtils;
import org.springframework.web.filter.OncePerRequestFilter;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.opensource.docgrid.global.common.response.ErrorResponse;
+import com.opensource.docgrid.global.exception.ErrorCode;
+
import io.jsonwebtoken.Claims;
import jakarta.servlet.FilterChain;
import jakarta.servlet.ServletException;
@@ -17,6 +23,12 @@
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
+/**
+ * HTTP JWT를 인증하고 요청별 현재 권한을 SecurityContext에 넣는다.
+ *
+ *
관리자 요청은 Redis 권한 캐시를 신뢰하지 않고 primary를 직접 확인한다.
+ * 확인할 수 없으면 이전 ADMIN 권한으로 진행시키지 않고 503을 반환한다.
+ */
@Slf4j
@RequiredArgsConstructor
public class JwtAuthenticationFilter extends OncePerRequestFilter {
@@ -24,6 +36,8 @@ public class JwtAuthenticationFilter extends OncePerRequestFilter {
private final JwtProvider jwtProvider;
private final TokenBlacklistService tokenBlacklistService;
private final RoleAuthorityService roleAuthorityService;
+ private final ObjectMapper objectMapper;
+ private final RequestMatcher adminRequests;
@Override
protected void doFilterInternal(HttpServletRequest request,
@@ -36,7 +50,25 @@ protected void doFilterInternal(HttpServletRequest request,
if (claims != null && !isBlacklisted(claims.get("jti", String.class))) {
Long userId = claims.get("userId", Long.class);
String email = claims.getSubject();
- List roles = roleAuthorityService.getRoles(userId);
+ // 1. 관리자 경로는 매번 primary에서 검증하고 일반 경로만 Redis 역할 캐시를 사용한다.
+ boolean adminPath = adminRequests.matches(request);
+ List roles;
+ try {
+ roles = adminPath ? roleAuthorityService.getRolesForAdmin(userId)
+ : roleAuthorityService.getRoles(userId);
+ } catch (RuntimeException e) {
+ if (!adminPath) {
+ throw e;
+ }
+ // 2. primary 확인이 불가능하면 캐시된 ADMIN으로 통과시키지 않는다.
+ log.error("관리자 권한 primary 검증 실패: {}", e.getClass().getSimpleName());
+ ErrorCode errorCode = ErrorCode.ADMIN_ROLE_UNAVAILABLE;
+ response.setStatus(errorCode.getHttpStatus().value());
+ response.setContentType(MediaType.APPLICATION_JSON_VALUE);
+ response.setCharacterEncoding("UTF-8");
+ objectMapper.writeValue(response.getWriter(), ErrorResponse.of(errorCode, request));
+ return;
+ }
List authorities = roles.stream()
.map(role -> new SimpleGrantedAuthority("ROLE_" + role))
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/auth/jwt/RoleAuthorityService.java b/backend/src/main/java/com/opensource/docgrid/domain/auth/jwt/RoleAuthorityService.java
index 758c0453..e40f50ee 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/auth/jwt/RoleAuthorityService.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/auth/jwt/RoleAuthorityService.java
@@ -5,22 +5,25 @@
import java.util.List;
import org.springframework.data.redis.core.StringRedisTemplate;
+import org.springframework.data.redis.core.script.DefaultRedisScript;
+import org.springframework.data.redis.core.script.RedisScript;
import org.springframework.stereotype.Component;
+import com.opensource.docgrid.domain.auth.service.query.PrimaryRoleQueryService;
import com.opensource.docgrid.domain.user.repository.UserRoleRepository;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
/**
- * 인가(hasRole) 판단에 쓰는 사용자 role을 JWT가 아니라 DB에서 매 요청 조회한다.
+ * 인가(hasRole) 판단에 쓰는 사용자 role을 JWT에 고정하지 않고 DB를 기준으로 관리한다.
*
* JWT에 role을 박제하면 관리자가 role을 부여/회수해도 재로그인 전까지 반영되지 않는다.
- * DB 조회 부하를 줄이기 위해 Redis에 짧은 TTL로 캐싱하고, role 변경 시 즉시 무효화한다.
+ * 일반 요청의 DB 조회 부하를 줄이기 위해 Redis에 짧은 TTL로 캐싱한다. HTTP 관리자
+ * 인가는 캐시를 우회해 primary를 매번 검증하며, role 변경 시 캐시 세대를 올린다.
*
- *
이 서비스는 인증 필터(모든 요청)의 critical path에 있으므로, {@code TokenBlacklistService}와
- * 동일하게 Redis 장애 시 예외를 전파하지 않고 DB 조회로 폴백한다 — Redis가 죽었다고 전체 API가
- * 막히면 안 된다.
+ *
일반 요청은 Redis 장애 시 기존과 같이 DB 조회로 폴백한다. 관리자 HTTP 요청은
+ * 가용성보다 최신 권한을 우선하며 primary 검증 실패를 호출자에게 전파한다.
*/
@Slf4j
@Component
@@ -28,10 +31,23 @@
public class RoleAuthorityService {
private static final String KEY_PREFIX = "auth:roles:";
+ private static final String EPOCH_PREFIX = "auth:roles:epoch:";
private static final Duration TTL = Duration.ofSeconds(30);
+ private static final RedisScript INVALIDATE = new DefaultRedisScript<>("""
+ redis.call('INCR', KEYS[2])
+ redis.call('DEL', KEYS[1])
+ return 1
+ """, Long.class);
+ private static final RedisScript CACHE_IF_UNCHANGED = new DefaultRedisScript<>("""
+ local current = redis.call('GET', KEYS[2]) or '0'
+ if current ~= ARGV[1] then return 0 end
+ redis.call('SET', KEYS[1], ARGV[2], 'EX', ARGV[3])
+ return 1
+ """, Long.class);
private final StringRedisTemplate redisTemplate;
private final UserRoleRepository userRoleRepository;
+ private final PrimaryRoleQueryService primaryRoleQueryService;
public List getRoles(Long userId) {
// 1. 먼저 Redis 캐시를 확인한다 — 대부분의 요청은 여기서 끝나 DB 부하를 줄인다.
@@ -40,17 +56,30 @@ public List getRoles(Long userId) {
return cached.isBlank() ? List.of() : Arrays.asList(cached.split(","));
}
- // 2. 캐시 미스면 DB에서 최신 role을 조회한다(source of truth).
+ // 2. 조회 전 세대를 기억해 회수·부여가 DB 조회와 캐시 저장 사이에 끼어드는지 확인한다.
+ String epoch = readEpoch(userId);
List roles = userRoleRepository.findRoleCodesByUserId(userId);
- // 3. 다음 요청부터는 캐시로 처리되도록 짧은 TTL로 저장해둔다.
- writeCache(userId, roles);
+ // 3. 세대가 바뀌면 조회 결과를 반환하거나 재캐시하지 않는다.
+ if (epoch != null && !writeCacheIfUnchanged(userId, epoch, roles)) {
+ return List.of();
+ }
return roles;
}
+ public List getRolesForAdmin(Long userId) {
+ // 관리자 인가는 Redis 상태와 복제 지연에 관계없이 현재 primary만 신뢰한다.
+ return primaryRoleQueryService.findCurrentRoles(userId);
+ }
+
public void invalidate(Long userId) {
try {
- redisTemplate.delete(KEY_PREFIX + userId);
+ // 세대 증가와 삭제가 원자적이어야 이전 DB 읽기가 삭제 뒤 캐시를 부활시키지 못한다.
+ Long invalidated = redisTemplate.execute(INVALIDATE,
+ List.of(KEY_PREFIX + userId, EPOCH_PREFIX + userId));
+ if (!Long.valueOf(1).equals(invalidated)) {
+ log.error("Redis role 캐시 무효화 결과를 확인할 수 없습니다. userId={}", userId);
+ }
} catch (Exception e) {
log.error("Redis role 캐시 무효화 실패, userId={}: {}", userId, e.getMessage());
}
@@ -65,11 +94,25 @@ private String readCache(Long userId) {
}
}
- private void writeCache(Long userId, List roles) {
+ private String readEpoch(Long userId) {
+ try {
+ String epoch = redisTemplate.opsForValue().get(EPOCH_PREFIX + userId);
+ return epoch == null ? "0" : epoch;
+ } catch (Exception e) {
+ log.error("Redis role 캐시 세대 조회 실패, DB로 폴백합니다. userId={}: {}", userId, e.getMessage());
+ return null;
+ }
+ }
+
+ private boolean writeCacheIfUnchanged(Long userId, String epoch, List roles) {
try {
- redisTemplate.opsForValue().set(KEY_PREFIX + userId, String.join(",", roles), TTL);
+ Long saved = redisTemplate.execute(CACHE_IF_UNCHANGED,
+ List.of(KEY_PREFIX + userId, EPOCH_PREFIX + userId), epoch,
+ String.join(",", roles), String.valueOf(TTL.toSeconds()));
+ return Long.valueOf(1).equals(saved);
} catch (Exception e) {
- log.error("Redis role 캐시 저장 실패, userId={}: {}", userId, e.getMessage());
+ log.error("Redis role 캐시 저장 실패, DB 조회 결과를 사용합니다. userId={}: {}", userId, e.getMessage());
+ return true;
}
}
}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/auth/service/query/PrimaryRoleQueryService.java b/backend/src/main/java/com/opensource/docgrid/domain/auth/service/query/PrimaryRoleQueryService.java
new file mode 100644
index 00000000..3d72eab6
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/auth/service/query/PrimaryRoleQueryService.java
@@ -0,0 +1,39 @@
+package com.opensource.docgrid.domain.auth.service.query;
+
+import java.util.List;
+
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Propagation;
+import org.springframework.transaction.annotation.Transactional;
+
+import com.opensource.docgrid.domain.user.repository.UserRoleRepository;
+
+import jakarta.persistence.EntityManager;
+import lombok.RequiredArgsConstructor;
+
+/**
+ * HTTP 관리자 인가에 필요한 역할을 Redis나 standby 없이 현재 primary에서 확인한다.
+ *
+ * 호출자의 read-only 트랜잭션을 상속하지 않으며, OpenProxy가 잘못 라우팅하면
+ * 권한을 추정하지 않고 요청을 실패시킨다. 일반 API·WebSocket의 역할 캐시는 담당하지 않는다.
+ */
+@Service
+@RequiredArgsConstructor
+public class PrimaryRoleQueryService {
+
+ private final EntityManager entityManager;
+ private final UserRoleRepository userRoleRepository;
+
+ @Transactional(propagation = Propagation.REQUIRES_NEW)
+ public List findCurrentRoles(Long userId) {
+ // 1. 명시적 read-write 트랜잭션 안에서 OpenProxy의 실제 도착 역할을 확인한다.
+ Object inRecovery = entityManager.createNativeQuery("SELECT pg_is_in_recovery()")
+ .getSingleResult();
+ if (!Boolean.FALSE.equals(inRecovery)) {
+ throw new IllegalStateException("Primary role verification is unavailable");
+ }
+
+ // 2. 같은 트랜잭션의 최신 primary에서 역할을 조회한다.
+ return userRoleRepository.findRoleCodesByUserId(userId);
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/user/service/command/UserRoleCommandService.java b/backend/src/main/java/com/opensource/docgrid/domain/user/service/command/UserRoleCommandService.java
index 7ec5262f..ef2672fb 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/user/service/command/UserRoleCommandService.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/user/service/command/UserRoleCommandService.java
@@ -81,8 +81,8 @@ public UserRoleResponse revokeRole(Long targetUserId, String roleCode) {
return UserRoleResponse.of(targetUser, roles);
}
- // DB 커밋 전에 캐시를 지우면, 커밋 직전 시점에 캐시 미스가 난 다른 요청이 아직 커밋 안 된(옛날) role을
- // 다시 캐시에 채워 넣을 수 있다. 그래서 무효화는 반드시 트랜잭션 커밋 이후로 미룬다.
+ // DB 커밋 전에 캐시를 지우면, 다른 요청이 아직 커밋 안 된 역할을 다시 읽을 수 있다.
+ // 커밋 후 무효화하고, 조회·저장 사이에 끼어드는 요청은 Redis 세대 비교로 재캐시를 막는다.
// 트랜잭션 밖에서 호출되는 경우(예: 단위 테스트)는 즉시 무효화한다.
private void invalidateAfterCommit(Long userId) {
if (TransactionSynchronizationManager.isSynchronizationActive()) {
diff --git a/backend/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java b/backend/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java
index d2c10427..dfc8ed5d 100644
--- a/backend/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java
+++ b/backend/src/main/java/com/opensource/docgrid/global/config/SecurityConfig.java
@@ -10,9 +10,12 @@
import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder;
import org.springframework.security.crypto.password.PasswordEncoder;
import org.springframework.security.web.SecurityFilterChain;
+import org.springframework.security.web.servlet.util.matcher.PathPatternRequestMatcher;
+import org.springframework.security.web.util.matcher.RequestMatcher;
import org.springframework.security.web.authentication.UsernamePasswordAuthenticationFilter;
import org.springframework.web.cors.CorsConfigurationSource;
+import com.fasterxml.jackson.databind.ObjectMapper;
import com.opensource.docgrid.domain.auth.jwt.JwtAuthenticationFilter;
import com.opensource.docgrid.domain.auth.jwt.JwtProvider;
import com.opensource.docgrid.domain.auth.jwt.RoleAuthorityService;
@@ -42,6 +45,7 @@ public class SecurityConfig {
private final McpAccessTokenCommandService mcpAccessTokenCommandService;
private final RestAuthenticationEntryPoint restAuthenticationEntryPoint;
private final RestAccessDeniedHandler restAccessDeniedHandler;
+ private final ObjectMapper objectMapper;
/**
* MCP({@code /mcp})와 웹 API({@code /mcp/tokens} 포함)를 포함한 전체 보안 필터 체인을 구성한다.
@@ -54,6 +58,8 @@ public class SecurityConfig {
@Bean
@Order(2)
public SecurityFilterChain filterChain(HttpSecurity http) throws Exception {
+ // 인가 규칙과 JWT 필터가 동일한 관리자 경로 판정을 사용해야 캐시 우회 틈이 없다.
+ RequestMatcher adminRequests = PathPatternRequestMatcher.withDefaults().matcher("/admin/**");
http
.cors(cors -> cors.configurationSource(corsConfigurationSource))
.csrf(AbstractHttpConfigurer::disable)
@@ -69,7 +75,7 @@ public SecurityFilterChain filterChain(HttpSecurity http) throws Exception {
// 2. STOMP CONNECT — StompAuthChannelInterceptor가 JWT 검증
// 3. STOMP SUBSCRIBE·SEND — StompDestinationAuthorizationInterceptor가 목적지별 권한 검증
.requestMatchers("/ws/**").permitAll()
- .requestMatchers("/admin/**").hasRole("ADMIN")
+ .requestMatchers(adminRequests).hasRole("ADMIN")
.anyRequest().authenticated()
)
// 인증 실패와 권한 부족을 상태 코드로 구분하고, 본문 없는 기본 응답 대신 공통 ErrorResponse를 준다.
@@ -82,7 +88,8 @@ public SecurityFilterChain filterChain(HttpSecurity http) throws Exception {
* "A를 B보다 앞자리에 꽂아라"는 뜻이라 위치 기준점(앵커)으로만 재사용한다 — 이 필터 앞에 꽂아야
* 두 인증 필터가 authorizeHttpRequests의 최종 인가 판정보다 먼저 실행돼 SecurityContext를 채울 수 있다.
*/
- .addFilterBefore(new JwtAuthenticationFilter(jwtProvider, tokenBlacklistService, roleAuthorityService), UsernamePasswordAuthenticationFilter.class)
+ .addFilterBefore(new JwtAuthenticationFilter(jwtProvider, tokenBlacklistService,
+ roleAuthorityService, objectMapper, adminRequests), UsernamePasswordAuthenticationFilter.class)
.addFilterBefore(new McpApiKeyAuthFilter(mcpAccessTokenCommandService), UsernamePasswordAuthenticationFilter.class);
return http.build();
}
diff --git a/backend/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java b/backend/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java
index 884e9de2..5d5e42d3 100644
--- a/backend/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java
+++ b/backend/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java
@@ -40,6 +40,7 @@ public enum ErrorCode {
PERMISSION_DENIED(HttpStatus.FORBIDDEN, "ROLE-002", "접근 권한이 없습니다."),
ROLE_ALREADY_ASSIGNED(HttpStatus.CONFLICT, "ROLE-003", "이미 부여된 역할입니다."),
ROLE_NOT_ASSIGNED(HttpStatus.NOT_FOUND, "ROLE-004", "부여되지 않은 역할입니다."),
+ ADMIN_ROLE_UNAVAILABLE(HttpStatus.SERVICE_UNAVAILABLE, "ROLE-005", "관리자 권한을 확인할 수 없습니다."),
// COLLECTION
COLLECTION_NOT_FOUND(HttpStatus.NOT_FOUND, "COLLECTION-001", "컬렉션을 찾을 수 없습니다."),
diff --git a/backend/src/test/java/com/opensource/docgrid/domain/auth/jwt/JwtAuthenticationFilterTest.java b/backend/src/test/java/com/opensource/docgrid/domain/auth/jwt/JwtAuthenticationFilterTest.java
index bf3b8d99..25a3aa41 100644
--- a/backend/src/test/java/com/opensource/docgrid/domain/auth/jwt/JwtAuthenticationFilterTest.java
+++ b/backend/src/test/java/com/opensource/docgrid/domain/auth/jwt/JwtAuthenticationFilterTest.java
@@ -3,6 +3,7 @@
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.BDDMockito.given;
+import static org.mockito.BDDMockito.then;
import java.util.List;
@@ -17,6 +18,9 @@
import org.springframework.mock.web.MockHttpServletRequest;
import org.springframework.mock.web.MockHttpServletResponse;
import org.springframework.security.core.context.SecurityContextHolder;
+import org.springframework.security.web.servlet.util.matcher.PathPatternRequestMatcher;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
@ExtendWith(MockitoExtension.class)
@DisplayName("JwtAuthenticationFilter 단위 테스트")
@@ -36,7 +40,9 @@ class JwtAuthenticationFilterTest {
@BeforeEach
void setUp() {
jwtProvider = new JwtProvider(TEST_SECRET, 3600L);
- filter = new JwtAuthenticationFilter(jwtProvider, tokenBlacklistService, roleAuthorityService);
+ filter = new JwtAuthenticationFilter(jwtProvider, tokenBlacklistService,
+ roleAuthorityService, new ObjectMapper().findAndRegisterModules(),
+ PathPatternRequestMatcher.withDefaults().matcher("/admin/**"));
SecurityContextHolder.clearContext();
}
@@ -80,9 +86,49 @@ void doFilter_authenticates_whenBlacklistCheckFails() throws Exception {
assertThat(SecurityContextHolder.getContext().getAuthentication()).isNotNull();
}
+ @Test
+ @DisplayName("관리자 요청은 Redis의 오래된 ADMIN 대신 primary의 현재 역할을 사용한다")
+ void doFilter_adminRequestUsesCurrentPrimaryRoles() throws Exception {
+ String token = jwtProvider.generateToken(1L, "user@test.com");
+ given(tokenBlacklistService.isBlacklisted(anyString())).willReturn(false);
+ given(roleAuthorityService.getRolesForAdmin(1L)).willReturn(List.of("USER"));
+
+ filter.doFilter(adminRequestWithToken(token), new MockHttpServletResponse(), new MockFilterChain());
+
+ assertThat(SecurityContextHolder.getContext().getAuthentication().getAuthorities())
+ .extracting("authority").containsExactly("ROLE_USER");
+ then(roleAuthorityService).should().getRolesForAdmin(1L);
+ then(roleAuthorityService).shouldHaveNoMoreInteractions();
+ }
+
+ @Test
+ @DisplayName("primary 역할 검증 실패 시 관리자 요청을 503으로 거부한다")
+ void doFilter_adminRequestFailsClosedWhenPrimaryUnavailable() throws Exception {
+ String token = jwtProvider.generateToken(1L, "user@test.com");
+ given(tokenBlacklistService.isBlacklisted(anyString())).willReturn(false);
+ given(roleAuthorityService.getRolesForAdmin(1L))
+ .willThrow(new IllegalStateException("primary unavailable"));
+ MockHttpServletResponse response = new MockHttpServletResponse();
+ MockFilterChain chain = new MockFilterChain();
+
+ filter.doFilter(adminRequestWithToken(token), response, chain);
+
+ assertThat(response.getStatus()).isEqualTo(503);
+ assertThat(response.getContentAsString()).contains("ROLE-005");
+ assertThat(chain.getRequest()).isNull();
+ assertThat(SecurityContextHolder.getContext().getAuthentication()).isNull();
+ }
+
private MockHttpServletRequest requestWithToken(String token) {
MockHttpServletRequest request = new MockHttpServletRequest();
request.addHeader("Authorization", "Bearer " + token);
return request;
}
+
+ private MockHttpServletRequest adminRequestWithToken(String token) {
+ MockHttpServletRequest request = requestWithToken(token);
+ request.setRequestURI("/admin/workers");
+ request.setServletPath("/admin/workers");
+ return request;
+ }
}
diff --git a/backend/src/test/java/com/opensource/docgrid/domain/auth/jwt/RoleAuthorityServiceTest.java b/backend/src/test/java/com/opensource/docgrid/domain/auth/jwt/RoleAuthorityServiceTest.java
index ae1f492a..0edd0a74 100644
--- a/backend/src/test/java/com/opensource/docgrid/domain/auth/jwt/RoleAuthorityServiceTest.java
+++ b/backend/src/test/java/com/opensource/docgrid/domain/auth/jwt/RoleAuthorityServiceTest.java
@@ -3,6 +3,8 @@
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.BDDMockito.given;
import static org.mockito.BDDMockito.then;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyList;
import java.util.List;
@@ -14,7 +16,9 @@
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.data.redis.core.ValueOperations;
+import org.springframework.data.redis.core.script.RedisScript;
+import com.opensource.docgrid.domain.auth.service.query.PrimaryRoleQueryService;
import com.opensource.docgrid.domain.user.repository.UserRoleRepository;
/**
@@ -28,6 +32,7 @@ class RoleAuthorityServiceTest {
@Mock private StringRedisTemplate redisTemplate;
@Mock private ValueOperations valueOperations;
@Mock private UserRoleRepository userRoleRepository;
+ @Mock private PrimaryRoleQueryService primaryRoleQueryService;
@Test
@DisplayName("캐시 히트: Redis에 값이 있으면 DB를 조회하지 않는다")
@@ -46,12 +51,18 @@ void getRoles_returnsCachedRoles_whenCacheHit() {
void getRoles_fetchesFromDbAndCaches_whenCacheMiss() {
given(redisTemplate.opsForValue()).willReturn(valueOperations);
given(valueOperations.get("auth:roles:1")).willReturn(null);
+ given(valueOperations.get("auth:roles:epoch:1")).willReturn(null);
given(userRoleRepository.findRoleCodesByUserId(1L)).willReturn(List.of("USER"));
+ given(redisTemplate.execute(any(RedisScript.class), anyList(), any(), any(), any()))
+ .willReturn(1L);
List roles = roleAuthorityService.getRoles(1L);
assertThat(roles).containsExactly("USER");
- then(valueOperations).should().set("auth:roles:1", "USER", java.time.Duration.ofSeconds(30));
+ then(redisTemplate).should().execute(any(RedisScript.class),
+ org.mockito.ArgumentMatchers.eq(List.of("auth:roles:1", "auth:roles:epoch:1")),
+ org.mockito.ArgumentMatchers.eq("0"), org.mockito.ArgumentMatchers.eq("USER"),
+ org.mockito.ArgumentMatchers.eq("30"));
}
@Test
@@ -59,7 +70,8 @@ void getRoles_fetchesFromDbAndCaches_whenCacheMiss() {
void invalidate_removesCacheKey() {
roleAuthorityService.invalidate(1L);
- then(redisTemplate).should().delete("auth:roles:1");
+ then(redisTemplate).should().execute(any(RedisScript.class),
+ org.mockito.ArgumentMatchers.eq(List.of("auth:roles:1", "auth:roles:epoch:1")));
}
@Test
@@ -72,4 +84,29 @@ void getRoles_fallsBackToDb_whenRedisReadFails() {
assertThat(roles).containsExactly("USER");
}
+
+ @Test
+ @DisplayName("DB 조회 중 권한이 바뀌면 오래된 역할을 반환하거나 재캐시하지 않는다")
+ void getRoles_discardsOldRoles_whenEpochChangesDuringDbRead() {
+ given(redisTemplate.opsForValue()).willReturn(valueOperations);
+ given(valueOperations.get("auth:roles:1")).willReturn(null);
+ given(valueOperations.get("auth:roles:epoch:1")).willReturn("0");
+ given(userRoleRepository.findRoleCodesByUserId(1L)).willAnswer(ignored -> {
+ roleAuthorityService.invalidate(1L);
+ return List.of("ADMIN");
+ });
+ given(redisTemplate.execute(any(RedisScript.class), anyList(), any(), any(), any()))
+ .willReturn(0L);
+
+ assertThat(roleAuthorityService.getRoles(1L)).isEmpty();
+ }
+
+ @Test
+ @DisplayName("관리자 요청은 Redis 역할 캐시를 읽지 않고 primary 조회 결과를 쓴다")
+ void getRolesForAdmin_usesPrimaryWithoutRedis() {
+ given(primaryRoleQueryService.findCurrentRoles(1L)).willReturn(List.of("USER"));
+
+ assertThat(roleAuthorityService.getRolesForAdmin(1L)).containsExactly("USER");
+ then(redisTemplate).shouldHaveNoInteractions();
+ }
}
diff --git a/backend/src/test/java/com/opensource/docgrid/domain/auth/service/query/PrimaryRoleQueryServiceTest.java b/backend/src/test/java/com/opensource/docgrid/domain/auth/service/query/PrimaryRoleQueryServiceTest.java
new file mode 100644
index 00000000..4ee93cc2
--- /dev/null
+++ b/backend/src/test/java/com/opensource/docgrid/domain/auth/service/query/PrimaryRoleQueryServiceTest.java
@@ -0,0 +1,69 @@
+package com.opensource.docgrid.domain.auth.service.query;
+
+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.then;
+
+import java.util.List;
+
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.InjectMocks;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+import org.springframework.transaction.annotation.Propagation;
+import org.springframework.transaction.annotation.Transactional;
+
+import com.opensource.docgrid.domain.user.repository.UserRoleRepository;
+
+import jakarta.persistence.EntityManager;
+import jakarta.persistence.Query;
+
+/**
+ * 관리자 역할 조회가 primary 확인 전에는 역할 SQL을 실행하지 않는 경계를 검증한다.
+ */
+@ExtendWith(MockitoExtension.class)
+@DisplayName("PrimaryRoleQueryService 단위 테스트")
+class PrimaryRoleQueryServiceTest {
+
+ @InjectMocks private PrimaryRoleQueryService service;
+ @Mock private EntityManager entityManager;
+ @Mock private Query query;
+ @Mock private UserRoleRepository userRoleRepository;
+
+ @Test
+ @DisplayName("primary에서만 현재 역할을 반환한다")
+ void findCurrentRoles_readsRolesOnPrimary() {
+ given(entityManager.createNativeQuery("SELECT pg_is_in_recovery()"))
+ .willReturn(query);
+ given(query.getSingleResult()).willReturn(false);
+ given(userRoleRepository.findRoleCodesByUserId(1L)).willReturn(List.of("USER"));
+
+ assertThat(service.findCurrentRoles(1L)).containsExactly("USER");
+ }
+
+ @Test
+ @DisplayName("standby라면 오래된 ADMIN을 조회하지 않고 실패한다")
+ void findCurrentRoles_rejectsStandby() {
+ given(entityManager.createNativeQuery("SELECT pg_is_in_recovery()"))
+ .willReturn(query);
+ given(query.getSingleResult()).willReturn(true);
+
+ assertThatThrownBy(() -> service.findCurrentRoles(1L))
+ .isInstanceOf(IllegalStateException.class);
+ then(userRoleRepository).shouldHaveNoInteractions();
+ }
+
+ @Test
+ @DisplayName("호출자의 read-only 경계를 상속하지 않는 새 트랜잭션을 요구한다")
+ void findCurrentRoles_startsIndependentReadWriteTransaction() throws NoSuchMethodException {
+ Transactional transaction = PrimaryRoleQueryService.class
+ .getMethod("findCurrentRoles", Long.class).getAnnotation(Transactional.class);
+
+ assertThat(transaction).isNotNull();
+ assertThat(transaction.readOnly()).isFalse();
+ assertThat(transaction.propagation()).isEqualTo(Propagation.REQUIRES_NEW);
+ }
+}
diff --git a/backend/src/test/java/com/opensource/docgrid/domain/user/controller/AdminUserControllerTest.java b/backend/src/test/java/com/opensource/docgrid/domain/user/controller/AdminUserControllerTest.java
index 4f785cd9..fb15b9ea 100644
--- a/backend/src/test/java/com/opensource/docgrid/domain/user/controller/AdminUserControllerTest.java
+++ b/backend/src/test/java/com/opensource/docgrid/domain/user/controller/AdminUserControllerTest.java
@@ -1,6 +1,7 @@
package com.opensource.docgrid.domain.user.controller;
import static org.mockito.BDDMockito.given;
+import static org.mockito.Mockito.mock;
import static org.springframework.security.test.web.servlet.request.SecurityMockMvcRequestPostProcessors.user;
import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.delete;
import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.get;
@@ -40,6 +41,8 @@
import com.opensource.docgrid.global.exception.DocGridException;
import com.opensource.docgrid.global.exception.ErrorCode;
+import io.jsonwebtoken.Claims;
+
/**
* 관리자 사용자 목록 API의 필터·Pagination·민감 정보 비노출과 ADMIN Security 계약을 검증한다.
*/
@@ -102,6 +105,50 @@ void getUsers_returnsForbidden_withoutAdminRole() throws Exception {
mockMvc.perform(get(USERS_URL)).andExpect(status().isUnauthorized());
}
+ @Test
+ @DisplayName("JWT의 관리자는 primary에서 ADMIN이 회수됐다면 캐시와 관계없이 403이다")
+ void getUsers_returnsForbidden_whenCurrentPrimaryRoleIsNotAdmin() throws Exception {
+ mockAdminJwt();
+ given(roleAuthorityService.getRolesForAdmin(10L)).willReturn(List.of("USER"));
+
+ mockMvc.perform(get(USERS_URL).header("Authorization", "Bearer current-token"))
+ .andExpect(status().isForbidden())
+ .andExpect(jsonPath("$.code").value("ROLE-002"));
+ }
+
+ @Test
+ @DisplayName("JWT의 ADMIN이 primary에 남아 있으면 관리자 요청을 허용한다")
+ void getUsers_allowsAdmin_whenCurrentPrimaryRoleIsAdmin() throws Exception {
+ mockAdminJwt();
+ given(roleAuthorityService.getRolesForAdmin(10L)).willReturn(List.of("ADMIN"));
+ given(adminUserQueryService.getUsers(null, null, null, 0, 20))
+ .willReturn(new PageResponse<>(List.of(), 0, 20, 0, 0, true, true));
+
+ mockMvc.perform(get(USERS_URL).header("Authorization", "Bearer current-token"))
+ .andExpect(status().isOk());
+ }
+
+ @Test
+ @DisplayName("primary 권한 확인에 실패하면 JWT 관리자 요청은 503이다")
+ void getUsers_returnsUnavailable_whenPrimaryCannotVerifyRole() throws Exception {
+ mockAdminJwt();
+ given(roleAuthorityService.getRolesForAdmin(10L))
+ .willThrow(new IllegalStateException("primary unavailable"));
+
+ mockMvc.perform(get(USERS_URL).header("Authorization", "Bearer current-token"))
+ .andExpect(status().isServiceUnavailable())
+ .andExpect(jsonPath("$.code").value("ROLE-005"));
+ }
+
+ private void mockAdminJwt() {
+ Claims claims = mock(Claims.class);
+ given(jwtProvider.getClaimsIfValid("current-token")).willReturn(claims);
+ given(claims.get("jti", String.class)).willReturn("current-jti");
+ given(claims.get("userId", Long.class)).willReturn(10L);
+ given(claims.getSubject()).willReturn("admin@example.com");
+ given(tokenBlacklistService.isBlacklisted("current-jti")).willReturn(false);
+ }
+
@Test
@DisplayName("페이지 입력 범위를 벗어나면 400을 반환한다")
void getUsers_returnsBadRequest_whenPageInputIsInvalid() throws Exception {
diff --git a/backend/src/test/java/com/opensource/docgrid/opensql/OpenSqlProxyJpaIntegrationTest.java b/backend/src/test/java/com/opensource/docgrid/opensql/OpenSqlProxyJpaIntegrationTest.java
index 12b9382a..959b34d9 100644
--- a/backend/src/test/java/com/opensource/docgrid/opensql/OpenSqlProxyJpaIntegrationTest.java
+++ b/backend/src/test/java/com/opensource/docgrid/opensql/OpenSqlProxyJpaIntegrationTest.java
@@ -4,6 +4,7 @@
import java.time.LocalDateTime;
+import com.opensource.docgrid.domain.auth.service.query.PrimaryRoleQueryService;
import com.opensource.docgrid.domain.failover.entity.FailoverEvent;
import com.opensource.docgrid.domain.failover.enums.FailoverEventType;
import com.opensource.docgrid.domain.failover.enums.FailoverStatus;
@@ -41,6 +42,9 @@ class OpenSqlProxyJpaIntegrationTest {
@Autowired
private EntityManager entityManager;
+ @Autowired
+ private PrimaryRoleQueryService primaryRoleQueryService;
+
@Test
@Timeout(30)
@DisplayName("명시적 JPA 트랜잭션은 리더에서 UPDATE와 엔티티 INSERT를 실행하고 롤백한다")
@@ -66,4 +70,12 @@ void jpaTransactionWritesToLeader() {
entityManager.flush();
assertThat(probe.getId()).isNotNull();
}
+
+ @Test
+ @Timeout(30)
+ @DisplayName("관리자 역할 조회의 독립 트랜잭션도 OpenProxy를 거쳐 primary에 도착한다")
+ void adminRoleQueryUsesPrimary() {
+ // 존재하지 않는 사용자만 조회해 운영 역할 데이터를 바꾸지 않는다.
+ assertThat(primaryRoleQueryService.findCurrentRoles(-1L)).isEmpty();
+ }
}
diff --git a/docs/test-results/opensql-admin-role-primary-consistency-20261001.md b/docs/test-results/opensql-admin-role-primary-consistency-20261001.md
new file mode 100644
index 00000000..efcebc38
--- /dev/null
+++ b/docs/test-results/opensql-admin-role-primary-consistency-20261001.md
@@ -0,0 +1,43 @@
+# OpenProxy 경유 관리자 권한 회수 정합성 검증
+
+## 목적과 변경 경계
+
+기존 재현 시험에서는 두 standby의 WAL 적용을 잠시 늦춘 뒤 ADMIN 역할을 회수했을 때, OpenProxy 경유 `/admin/**` 요청이 뒤처진 역할을 읽어 200으로 통과했다. 이 변경은 **HTTP 관리자 인가**에 한해 매 요청을 명시적 read-write 트랜잭션의 현재 primary에서 검증한다. 같은 연결에서 `pg_is_in_recovery() = false`를 확인하지 못하거나 DB 조회가 실패하면 캐시된 ADMIN을 사용하지 않고 503을 반환한다. 일반 API의 기존 Redis 캐시 경로와 WebSocket 세션 인가는 이번 변경 범위가 아니다.
+
+일반 역할 캐시에는 조회 전 epoch를 기억하고 Lua CAS로 저장하는 장치를 추가했다. 역할 변경 후에는 epoch 증가와 캐시 삭제를 하나의 Redis 스크립트로 실행한다. 이 장치는 **DB 조회와 캐시 저장 사이의 무효화 경쟁**을 막는다. 그러나 일반 경로가 무효화 이후 새 epoch를 읽고 지연된 standby에서 옛 권한을 읽는 경우까지 해결하지 않는다. 따라서 “모든 권한 조회가 최신”이라는 결론은 내리지 않는다. 관리자 HTTP 경로가 캐시와 standby를 우회하는 것이 보안 경계다.
+
+## 검증 환경과 절차
+
+- 실행: 2026-10-01 KST. 로컬의 격리된 앱 두 인스턴스에서 loopback SSH 터널로 GCP 3노드 OpenSQL에 접속했다. 한 인스턴스는 OpenProxy 두 주소, 다른 인스턴스는 primary 직접 경로를 사용했다. 임시 Redis는 별도 loopback 포트의 자동 삭제 컨테이너였다.
+- 코드: `fix/admin-role-revocation-consistency`의 `ac6bd86`까지. 실제 실행 JAR SHA-256은 `8317e9c5e1d17aa21e9fde941c379412cc39d6749c80b0d65eb73afd6928fe57`이다. 비밀이 없는 live preflight SHA-256은 `ba05ec8fd820ea2201c2448d88b482b461a3aa576423ed7e627fa63aa7c4564e`다.
+- 비밀번호는 승인된 SSH 스트림에서 테스트 프로세스 메모리로만 전달했다. 공개 문서·커밋·원장에는 값과 서버 주소를 넣지 않았다. 시험 전용 사용자 4개는 보호 타이머 아래 생성·삭제했다.
+- 사전 관문: 실제 GCP 계정·프로젝트, 1 primary/2 streaming standby, 설치된 보호 helper 해시, 두 standby의 정상 복제·승격 자격, ADMIN 기본 200, 역할 SQL의 primary 도착을 확인했다.
+- 장애 주입: 두 standby 각각에 독립 복구 타이머를 먼저 예약하고 `nofailover=true`와 120초 apply delay를 적용했다. 해당 구간의 primary 장애 시 자동 failover가 제한될 수 있으므로 이 상태를 지속 운영 구성으로 해석하면 안 된다. 회수 API가 200을 반환한 뒤 OpenProxy 요청과 직접 primary 대조 요청을 실행했다.
+- 정리: 지연을 즉시 해제하고 정상 복제를 확인한 뒤 시험 역할·사용자를 삭제하고 보호 타이머를 취소했다. 별도 read-only 재확인에서 시험 계정 0건, node2·3 모두 `streaming=true`, `armed=false`, `delay=0`, `nofailover=false`, WAL backlog 0이었다. 임시 Redis와 SSH 터널도 종료했다.
+
+동일 방식의 재실행 진입점은 `scripts/opensql/run_permission_replica_lag_local.py --credentials-stdin --apply-delay --expect-primary-admin`이다. `--credentials-stdin`은 승인된 SSH 표준출력에만 연결하고, 다른 환경값은 비공개 `.env`와 `OPENSQL_*` 변수로 제공한다. 이 명령은 실제 standby 정책을 일시 변경하므로 독립 복구 타이머와 운영 승인 없이 실행하지 않는다.
+
+## 관측값과 해석
+
+원장 run ID `2075ca4865f1`에서 16요청의 결과는 2xx 13건, 명시적 실패 3건, 결과 불명 0건이었다. 장애 구간은 UTC `2026-09-30T17:56:57.911Z`부터 `17:58:02.784Z`까지였다. 아래 SQL 수는 해당 단계 전후 `pg_stat_statements`의 역할 조회 증가량이며, node1은 primary, node2·3은 당시 standby였다.
+
+| 단계 | DB 역할 상태 | OpenProxy 관리자 HTTP | 역할 SQL 증가량 node1/node2/node3 | Redis ADMIN 캐시 | 판정 |
+| --- | --- | ---: | ---: | --- | --- |
+| R: 정상 라우팅 6회 | 세 노드 모두 최신 | 6회 모두 200 | 6 / 0 / 0 | 캐시 미사용 | 관리자 권한 조회가 primary에 고정됨 |
+| M: ADMIN 회수 직후 | primary에는 회수 반영, 두 standby에는 이전 ADMIN 유지 | **403** | **1 / 0 / 0** | 없음 | 뒤처진 standby가 있어도 권한 우회 없음 |
+| C2: 직접 primary 대조 | M과 동일 | 직접 경로 **403** | 1 / 0 / 0 | 없음 | 새 권한 기준 대조와 일치 |
+| C1: 복제 복구 후 회수 | 세 노드 모두 회수 반영 | **403** | 1 / 0 / 0 | 없음 | 정상 상태에서도 같은 정책 유지 |
+
+M 단계에서는 두 standby 모두 `armed=true`, `nofailover=true`, `delay=120000ms`였고, 각각 primary에는 없는 ADMIN 역할이 남아 있었다. 그런데 M 요청은 primary에서 역할 SQL이 1회 실행되고 403이었다. 따라서 이 403을 단순한 인증 실패나 standby 미지연으로 오인하지 않는다. 실행 분류는 `PRIMARY_AUTH_ENFORCED`였다.
+
+`scripts/opensql/ha_evidence.py verify --run-dir ` 결과 `complete=true`, 열린 장애 구간 0, 결과 불명 0이었다. 원본 `manifest.json`, `events.jsonl`, `requests.csv`, `permission-scenarios.json`, 앱 로그와 live preflight는 서버 식별자·내부 경로가 포함될 수 있어 비공개 실행 폴더에만 보관하며 Git에는 올리지 않는다.
+
+## 자동 회귀 검증과 한계
+
+- 관리자 Security filter 전체 경로: 정상 ADMIN 200, 회수 후 403, primary 조회 불가 시 503, 일반 API 기존 동작을 테스트했다.
+- 역할 캐시: Redis epoch 변화 시 이전 DB 조회 결과가 캐시에 다시 쓰이지 않도록 단위 테스트와 격리 Redis Lua 실행으로 확인했다.
+- `python3 -m unittest discover -s scripts/opensql -p 'test_permission_replica_lag.py'`: 6개 통과. 관련 Java 단위·MVC 테스트도 Gradle `test`에서 통과했다.
+- 전체 `./backend/gradlew -p backend test` 실행은 1,188개 중 117개가 실패했다. 이 실행 환경에서 기본 테스트 프로필이 기대하는 로컬 PostgreSQL 포트에 서버가 없어 Flyway의 JDBC 연결이 거절됐고, 같은 Spring context의 후속 시험도 로딩 실패로 이어졌다. 전체 테스트 통과로 표기하지 않는다.
+- `OpenSqlProxyJpaIntegrationTest`는 별도 외부 클러스터 태그로 일반 Gradle `test`에서 제외된다. 이번 실제 경로 증거는 위 HTTP 실행과 노드별 SQL 계측이다.
+- 503 경로는 로컬 보안 체인 테스트로 검증했다. 실제 primary 장애를 일으키는 503 실험은 이번에 하지 않았다.
+- WebSocket CONNECT·SUBSCRIBE·기존 구독의 역할 수명, 일반 API의 standby 지연 후 캐시 최신성, Redis 무효화 실패나 DB failover 중 이미 성공 응답된 회수의 RPO는 이 실험으로 증명되지 않는다. 별도 검증이 필요하다.
diff --git a/scripts/opensql/permission_replica_lag.py b/scripts/opensql/permission_replica_lag.py
index 3cea8dc8..39d54438 100644
--- a/scripts/opensql/permission_replica_lag.py
+++ b/scripts/opensql/permission_replica_lag.py
@@ -291,6 +291,17 @@ def classify_result(m_http: int, m_admin_cached: bool, m_delta: dict[str, int],
return "INVALID"
+def classify_primary_result(m_http: int, m_admin_cached: bool, m_delta: dict[str, int],
+ c1_http: int, c2_http: int) -> str:
+ # A denial without a measured primary role lookup does not prove the fix.
+ primary_lookup = (m_delta[NODES[0]] > 0 and
+ m_delta[NODES[1]] + m_delta[NODES[2]] == 0)
+ if primary_lookup and m_http == 403 and not m_admin_cached and \
+ c1_http == 403 and c2_http == 403:
+ return "PRIMARY_AUTH_ENFORCED"
+ return "INVALID"
+
+
def require_status(actual: int, expected: int, label: str) -> None:
if actual != expected:
raise InvalidExperiment(f"{label}: expected HTTP {expected}, received {actual}")
@@ -353,7 +364,8 @@ def run(args: argparse.Namespace) -> None:
separators=(",", ":")).encode() + b"\n"
config_hash = hashlib.sha256(snapshot).hexdigest()
run_dir = args.output / f"permission-{run_id}"
- ledger(run_dir, "init", "--scenario", "permission-replica-lag",
+ ledger(run_dir, "init", "--scenario", ("permission-primary-consistency"
+ if args.expect_primary_admin else "permission-replica-lag"),
"--run-id", run_id,
"--config-sha256", config_hash,
"--opensql-version", versions["opensql"],
@@ -362,6 +374,7 @@ def run(args: argparse.Namespace) -> None:
"--etcd-version", versions["etcd"])
write_private(run_dir / "live-preflight.json", snapshot)
evidence: dict[str, object] = {"run_id": run_id, "started_at": utc_now(),
+ "expect_primary_admin": args.expect_primary_admin,
"contract_sha256": contract["evidence_sha256"],
"live_preflight_sha256": config_hash,
"app_jar_sha256": args.app_jar_sha256,
@@ -393,8 +406,11 @@ def run(args: argparse.Namespace) -> None:
f"routing-{index}"), 200, "routing baseline")
route = deltas(baseline, role_stats())
evidence["R"] = {"role_sql_delta": route, "cache": cache(ids["m"])}
- if route[NODES[1]] + route[NODES[2]] == 0:
- raise InvalidExperiment("Role SQL did not reach standby; replay pause is prohibited")
+ if args.expect_primary_admin:
+ if route[NODES[0]] == 0 or route[NODES[1]] + route[NODES[2]] != 0:
+ raise InvalidExperiment("Admin role SQL did not reach only primary in baseline")
+ elif route[NODES[1]] + route[NODES[2]] == 0:
+ raise InvalidExperiment("Role SQL did not reach standby; apply delay is prohibited")
if not args.apply_delay:
raise InvalidExperiment("R recorded; pass --apply-delay only after the independent "
"reset timer and one-standby delay have been verified")
@@ -473,13 +489,18 @@ def run(args: argparse.Namespace) -> None:
time.sleep(1)
if any(roles(ids["c1"]).values()):
raise InvalidExperiment("C1 did not replicate before its fresh-read check")
+ c1_before = role_stats()
c1_http = http(run_dir, args.proxy_url, "/admin/workers", jwt["c1"], "fresh-C1")
+ c1_delta = deltas(c1_before, role_stats())
evidence["C1"] = {"http_status": c1_http, "cache": cache(ids["c1"]),
- "roles": roles(ids["c1"])}
+ "roles": roles(ids["c1"]), "role_sql_delta": c1_delta}
if c2_delta[NODES[0]] == 0 or c2_delta[NODES[1]] + c2_delta[NODES[2]] != 0:
raise InvalidExperiment("C2 did not prove direct-primary role lookup")
- evidence["result"] = classify_result(m_http, m_cache["admin"], m_delta,
- c1_http, c2_http)
+ if args.expect_primary_admin and (c1_delta[NODES[0]] == 0 or
+ c1_delta[NODES[1]] + c1_delta[NODES[2]] != 0):
+ raise InvalidExperiment("C1 did not prove primary-only role lookup")
+ classifier = classify_primary_result if args.expect_primary_admin else classify_result
+ evidence["result"] = classifier(m_http, m_cache["admin"], m_delta, c1_http, c2_http)
evidence["observation"] = evidence["result"]
except Exception as error:
evidence["result"] = "INVALID"
@@ -549,6 +570,8 @@ def main() -> int:
parser.add_argument("--app-jar-sha256", required=True)
parser.add_argument("--apply-delay", action="store_true",
help="Apply a temporary standby delay after independent reset verification")
+ parser.add_argument("--expect-primary-admin", action="store_true",
+ help="Require the fixed HTTP administrator path to read only primary")
args = parser.parse_args()
try:
if not re.fullmatch(r"[0-9a-f]{64}", args.app_jar_sha256):
diff --git a/scripts/opensql/run_permission_replica_lag_local.py b/scripts/opensql/run_permission_replica_lag_local.py
index 3462bf14..7290edcb 100644
--- a/scripts/opensql/run_permission_replica_lag_local.py
+++ b/scripts/opensql/run_permission_replica_lag_local.py
@@ -22,9 +22,9 @@
RUNNER = ROOT / "scripts/opensql/permission_replica_lag.py"
-def load_env(path: Path) -> dict[str, str]:
+def parse_env(contents: str) -> dict[str, str]:
values = {}
- for line in path.read_text().splitlines():
+ for line in contents.splitlines():
line = line.strip()
if not line or line.startswith("#"):
continue
@@ -35,6 +35,18 @@ def load_env(path: Path) -> dict[str, str]:
return values
+def load_env(path: Path) -> dict[str, str]:
+ return parse_env(path.read_text())
+
+
+def merge_credentials(values: dict[str, str], credentials: dict[str, str]) -> dict[str, str]:
+ # Map only the two approved DB passwords; no remote metadata enters app logs.
+ return values | {
+ "OPENSQL_APP_PASSWORD": credentials["DOCGRID_APP_PASSWORD"],
+ "OPENSQL_MIGRATION_PASSWORD": credentials["DOCGRID_MIGRATION_PASSWORD"],
+ }
+
+
def port_listening(port: int) -> bool:
try:
with socket.create_connection(("127.0.0.1", port), timeout=0.5):
@@ -94,18 +106,29 @@ def stop_app(process: subprocess.Popen) -> None:
def main() -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--env-file", type=Path, required=True)
+ parser.add_argument("--credentials-stdin", action="store_true",
+ help="Read remote DocGrid credentials through an encrypted SSH pipe")
parser.add_argument("--jar", type=Path, required=True)
parser.add_argument("--output", type=Path, required=True)
parser.add_argument("--redis-port", type=int, default=16379,
help="Port of a disposable, isolated loopback Redis instance")
parser.add_argument("--apply-delay", action="store_true",
help="Run the guarded standby apply-delay phase after local apps start")
+ parser.add_argument("--expect-primary-admin", action="store_true",
+ help="Require the fixed HTTP administrator path to read only primary")
args = parser.parse_args()
if args.output.exists() or any(port_listening(port) for port in (18080, 18081)):
print("Output directory or test app port already exists", file=sys.stderr)
return 2
# Non-secret route overrides may come from the process; the private file wins on conflicts.
values = os.environ | load_env(args.env_file)
+ if args.credentials_stdin:
+ # 1. The SSH pipe is consumed once and never copied to the evidence directory.
+ try:
+ values = merge_credentials(values, parse_env(sys.stdin.read()))
+ except (KeyError, ValueError):
+ print("Remote credential stream is incomplete", file=sys.stderr)
+ return 2
required = ("OPENSQL_APP_JDBC_URL", "OPENSQL_APP_DIRECT_JDBC_URL",
"OPENSQL_APP_USER", "OPENSQL_APP_PASSWORD", "JWT_SECRET")
if any(not values.get(key) for key in required):
@@ -133,6 +156,8 @@ def main() -> int:
]
if args.apply_delay:
runner_args.append("--apply-delay")
+ if args.expect_primary_admin:
+ runner_args.append("--expect-primary-admin")
result = subprocess.run(runner_args, cwd=ROOT, env=base_env, check=False)
return result.returncode
except (OSError, RuntimeError) as error:
diff --git a/scripts/opensql/test_permission_replica_lag.py b/scripts/opensql/test_permission_replica_lag.py
index 21ad624a..530566e9 100644
--- a/scripts/opensql/test_permission_replica_lag.py
+++ b/scripts/opensql/test_permission_replica_lag.py
@@ -13,11 +13,27 @@
raise RuntimeError("Permission experiment module cannot be loaded")
EXPERIMENT = importlib.util.module_from_spec(SPEC)
SPEC.loader.exec_module(EXPERIMENT)
+WRAPPER_SPEC = importlib.util.spec_from_file_location(
+ "run_permission_replica_lag_local", SCRIPT.with_name("run_permission_replica_lag_local.py"))
+if WRAPPER_SPEC is None or WRAPPER_SPEC.loader is None:
+ raise RuntimeError("Local experiment wrapper cannot be loaded")
+WRAPPER = importlib.util.module_from_spec(WRAPPER_SPEC)
+WRAPPER_SPEC.loader.exec_module(WRAPPER)
class PermissionReplicaLagTest(unittest.TestCase):
"""Keep a denied HTTP response distinct from a proven standby denial."""
+ def test_remote_credentials_are_mapped_in_memory(self):
+ """The SSH stream supplies only DB passwords and does not become evidence."""
+ credentials = WRAPPER.parse_env("DOCGRID_APP_PASSWORD=app-secret\n"
+ "DOCGRID_MIGRATION_PASSWORD=migration-secret\n")
+ merged = WRAPPER.merge_credentials({"JWT_SECRET": "local"}, credentials)
+ self.assertEqual("app-secret", merged["OPENSQL_APP_PASSWORD"])
+ self.assertEqual("migration-secret", merged["OPENSQL_MIGRATION_PASSWORD"])
+ self.assertEqual("local", merged["JWT_SECRET"])
+ self.assertNotIn("DOCGRID_APP_PASSWORD", merged)
+
def test_403_without_standby_query_is_invalid(self):
"""An authentication failure or primary read must not appear as a safe denial."""
no_query = dict.fromkeys(EXPERIMENT.NODES, 0)
@@ -38,6 +54,20 @@ def test_standby_read_separates_stale_and_denied(self):
self.assertEqual("INVALID", EXPERIMENT.classify_result(
200, True, standby_query, 200, 403))
+ def test_primary_mode_requires_denial_and_a_measured_primary_role_lookup(self):
+ """A 403 from authentication or a standby read is not a fixed-route proof."""
+ primary_query = {EXPERIMENT.NODES[0]: 1, EXPERIMENT.NODES[1]: 0,
+ EXPERIMENT.NODES[2]: 0}
+ standby_query = primary_query | {EXPERIMENT.NODES[0]: 0, EXPERIMENT.NODES[1]: 1}
+ self.assertEqual("PRIMARY_AUTH_ENFORCED", EXPERIMENT.classify_primary_result(
+ 403, False, primary_query, 403, 403))
+ self.assertEqual("INVALID", EXPERIMENT.classify_primary_result(
+ 403, False, standby_query, 403, 403))
+ self.assertEqual("INVALID", EXPERIMENT.classify_primary_result(
+ 200, True, primary_query, 403, 403))
+ self.assertEqual("INVALID", EXPERIMENT.classify_primary_result(
+ 403, True, primary_query, 403, 403))
+
def test_standby_status_requires_both_tag_and_timer(self):
"""A delay alone cannot authorize a revocation while replicas remain promotable."""
delayed = ("standby=true streaming=true armed=true delay=120000,configuration file "