diff --git a/.claude/launch.json b/.claude/launch.json new file mode 100644 index 00000000..51a03491 --- /dev/null +++ b/.claude/launch.json @@ -0,0 +1,11 @@ +{ + "version": "0.0.1", + "configurations": [ + { + "name": "backend", + "runtimeExecutable": "./backend/gradlew", + "runtimeArgs": ["-p", "backend", "bootRun", "--args=--spring.profiles.active=local"], + "port": 8080 + } + ] +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/config/RagExecutionConfig.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/config/RagExecutionConfig.java new file mode 100644 index 00000000..3fa0d212 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/config/RagExecutionConfig.java @@ -0,0 +1,78 @@ +package com.opensource.docgrid.domain.rag.config; + +import java.util.concurrent.Semaphore; +import java.util.concurrent.SynchronousQueue; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.concurrent.CustomizableThreadFactory; + +/** + * RAG 답변 생성을 동시에 몇 건까지 처리할지 제한하는 실행 자원을 구성한다 (#340). + * + *
{@code domain/worker/config/WorkerExecutionConfig}(이 코드베이스에서 유일했던 커스텀 + * 스레드풀 선례)를 그대로 본떴다 — {@code core=max}인 {@link ThreadPoolExecutor} + + * {@link SynchronousQueue}(큐잉 없음) + {@link ThreadPoolExecutor.AbortPolicy}(꽉 차면 즉시 + * 거부, 조용히 쌓아두지 않음)로 "설정된 동시성만 즉시 실행"을 보장한다. + * + *
{@code embedding_jobs}의 {@code WorkerExecutionSlotPool}(별도 클래스, 종료 플래그 + + * introspection 메서드 포함)까지는 필요 없다 — RAG는 별도 워커 등록/우아한 종료 조율이나 + * 대시보드 노출 요구가 없어서, 같은 안전 성질(로컬 슬롯을 먼저 확보한 뒤에만 DB claim을 + * 시도해 "claim은 됐는데 실행할 스레드가 없는" 상태를 만들지 않는 것)을 순수 + * {@link Semaphore}만으로 재현한다. + * + *
{@code indexing.worker.enabled}로 켜고 끌 수 있는 인덱싱 워커와 달리, RAG Worker는
+ * {@code RagSchedulingConfig}와 동일하게 조건 없이 항상 켜져 있어야 하는 검색 API 핵심
+ * 경로라 {@code @ConditionalOnProperty}를 붙이지 않는다.
+ */
+@Configuration
+public class RagExecutionConfig {
+
+ public static final String RAG_WORKER_JOB_EXECUTOR = "ragWorkerJobExecutor";
+ public static final String RAG_WORKER_SLOTS = "ragWorkerSlots";
+
+ /**
+ * 동시에 최대 {@code rag.worker.max-concurrency}건까지만 즉시 실행하는 무대기 Executor를 만든다.
+ */
+ @Bean(name = RAG_WORKER_JOB_EXECUTOR, destroyMethod = "shutdownNow")
+ public ThreadPoolExecutor ragWorkerJobExecutor(
+ @Value("${rag.worker.max-concurrency:2}") int maxConcurrency
+ ) {
+ validateMaxConcurrency(maxConcurrency);
+ return new ThreadPoolExecutor(
+ maxConcurrency,
+ maxConcurrency,
+ 0L,
+ TimeUnit.MILLISECONDS,
+ new SynchronousQueue<>(),
+ new CustomizableThreadFactory("rag-worker-job-"),
+ new ThreadPoolExecutor.AbortPolicy()
+ );
+ }
+
+ /**
+ * DB claim을 시도하기 전에 먼저 확보해야 하는 로컬 실행 슬롯. 공정 모드(fair)로 만들어
+ * 폴링 주기가 겹칠 때 대기가 한쪽으로 몰리지 않게 한다 — {@code WorkerExecutionSlotPool}의
+ * 선택과 동일하다.
+ */
+ @Bean(name = RAG_WORKER_SLOTS)
+ public Semaphore ragWorkerSlots(@Value("${rag.worker.max-concurrency:2}") int maxConcurrency) {
+ validateMaxConcurrency(maxConcurrency);
+ return new Semaphore(maxConcurrency, true);
+ }
+
+ /**
+ * {@code max-concurrency=0}(또는 음수)은 예외 없이 Worker를 영구 대기 상태로 만들 수 있어
+ * ({@link Semaphore}는 permit 0으로도 생성 자체는 허용) 기동 시점에 바로 실패시킨다.
+ */
+ private void validateMaxConcurrency(int maxConcurrency) {
+ if (maxConcurrency < 1) {
+ throw new IllegalArgumentException(
+ "rag.worker.max-concurrency는 1 이상이어야 합니다: " + maxConcurrency
+ );
+ }
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/entity/RagResponse.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/entity/RagResponse.java
index 73532eb3..3c6498c1 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/rag/entity/RagResponse.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/entity/RagResponse.java
@@ -1,5 +1,7 @@
package com.opensource.docgrid.domain.rag.entity;
+import java.time.LocalDateTime;
+
import com.opensource.docgrid.domain.search.entity.SearchQuery;
import com.opensource.docgrid.domain.search.enums.ResultStatus;
import com.opensource.docgrid.global.common.entity.BaseEntity;
@@ -88,6 +90,12 @@ public class RagResponse extends BaseEntity {
@Column(name = "error_message", columnDefinition = "TEXT")
private String errorMessage;
+ // 병렬 Worker가 이 job을 이미 집었는지 표시한다(#340). status만으로는 "대기 중"과 "누가 이미
+ // 처리 중"을 구분할 수 없어서(둘 다 PROCESSING) 별도로 둔다. RagResponseClaimService의 짧은
+ // claim 트랜잭션 안에서만 채워지며, 그 밖의 완료 확정 경로(조건부 UPDATE)는 이 컬럼을 건드리지 않는다.
+ @Column(name = "claimed_at")
+ private LocalDateTime claimedAt;
+
@Builder
public RagResponse(SearchQuery query, String answerText, String llmProvider, String llmModelName,
String promptText, Integer inputTokenCount, Integer outputTokenCount, Integer latencyMs,
@@ -103,4 +111,13 @@ public RagResponse(SearchQuery query, String answerText, String llmProvider, Str
this.status = status;
this.errorMessage = errorMessage;
}
+
+ /**
+ * 이 job을 지금 이 Worker가 처리하기 시작했다는 표시를 남긴다. {@code RagResponseClaimService}의
+ * 짧은 claim 트랜잭션 안에서만 호출되어야 한다 — dirty checking으로 반영되므로, 이 엔티티가
+ * detached된 뒤(다른 트랜잭션/스레드로 넘어간 뒤)에 호출하면 반영되지 않는다.
+ */
+ public void markClaimed(LocalDateTime claimedAt) {
+ this.claimedAt = claimedAt;
+ }
}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/repository/RagResponseRepository.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/repository/RagResponseRepository.java
index e25e2e08..db4defe6 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/rag/repository/RagResponseRepository.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/repository/RagResponseRepository.java
@@ -22,17 +22,49 @@
public interface RagResponseRepository extends JpaRepository {@code query}/{@code query.user}를 {@link EntityGraph}로 미리 fetch한다 — Worker가
- * 이 메서드로 job을 꺼낸 트랜잭션이 끝난 뒤(WebSocket push 시점)에
- * {@code job.getQuery().getUser().getEmail()}에 접근해도 두 연관관계 모두 LAZY라서 자칫
- * {@code LazyInitializationException}이 날 수 있는데, 미리 로딩해두면 그 문제가 없다.
+ * PROCESSING 중 아직 아무 Worker도 집지 않은(claim 안 된) 것 하나를 골라 행 잠금을 건다(#340).
+ * {@code claimed_at IS NULL} 조건이 "대기 중"과 "이미 처리 중"을 구분하는 유일한 신호다 —
+ * status만으로는 둘 다 PROCESSING이라 구분이 안 된다. {@code FOR UPDATE SKIP LOCKED}로
+ * 여러 Worker가 동시에 이 쿼리를 날려도 이미 잠긴 행은 건너뛰고 그다음 미잠금 행을 잡아온다.
+ * {@code created_at}이 같은 밀리초를 공유할 수 있는 동시 접수 상황을 대비해 {@code id}를
+ * 2차 정렬 기준으로 둔다. {@link RagResponseClaimService}가 이 메서드로 잠근 행을 같은 짧은
+ * 트랜잭션 안에서 즉시 {@link RagResponse#markClaimed}로 확정하고 커밋해, 락을 오래 들고
+ * 있지 않는다({@code embedding_jobs}의 claim 패턴과 동일).
+ */
+ @Query(value = """
+ SELECT * FROM rag_responses
+ WHERE status = 'PROCESSING' AND claimed_at IS NULL
+ ORDER BY created_at ASC, id ASC
+ LIMIT 1
+ FOR UPDATE SKIP LOCKED
+ """, nativeQuery = true)
+ Optional {@code job} 객체가 아니라 {@code jobId}만 받아 이 메서드 자신의 트랜잭션 안에서 다시
- * 조회하는 이유: RagJobWorker가 {@code findFirstByStatusOrderByCreatedAtAsc()}로 꺼낸
+ * 조회하는 이유: RagJobWorker가 claim 단계({@code RagResponseClaimService}, #340)에서 꺼낸
* job은 그 조회 시점에 트랜잭션이 끝나 detached 상태다. 원래(#218) 이 detached 인스턴스를
* 그대로 받아 필드만 바꾸면 dirty checking이 감지 못해 DB에 반영되지 않는 버그가 있었는데,
* 지금은 완료 처리 자체가 dirty checking에 의존하지 않는다({@link
@@ -229,7 +229,11 @@ public boolean processJob(Long jobId) {
}
// 답변과 citation 저장이 모두 끝난 동일 Transaction의 커밋 이후 성공 Counter를 기록한다.
applicationEventPublisher.publishEvent(new RagJobCompletionMetricEvent(Outcome.SUCCESS));
- log.info("[RAG] done queryId={} responseId={} latencyMs={}", queryId, job.getId(), result.latencyMs());
+ // promptTokens/answerTokens을 함께 남겨, 느린 job이 프롬프트를 읽느라(prefill) 오래 걸린 건지
+ // 답변을 쓰느라(decode) 오래 걸린 건지 로그만으로 구분할 수 있게 한다 — 병렬화(#340) 이후
+ // 요청당 작업량을 어느 쪽부터 줄여야 할지 판단하는 근거 자료.
+ log.info("[RAG] done queryId={} responseId={} latencyMs={} promptTokens={} answerTokens={}",
+ queryId, job.getId(), result.latencyMs(), result.inputTokenCount(), result.outputTokenCount());
return true;
}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/service/RagJobWorker.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/service/RagJobWorker.java
index a648a390..c164f145 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/rag/service/RagJobWorker.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/service/RagJobWorker.java
@@ -1,45 +1,130 @@
package com.opensource.docgrid.domain.rag.service;
import java.util.Optional;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.Semaphore;
+import java.util.concurrent.ThreadPoolExecutor;
+import org.springframework.beans.factory.annotation.Qualifier;
+import org.springframework.boot.context.event.ApplicationReadyEvent;
+import org.springframework.context.event.EventListener;
import org.springframework.dao.OptimisticLockingFailureException;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
+import com.opensource.docgrid.domain.rag.config.RagExecutionConfig;
import com.opensource.docgrid.domain.rag.controller.RagWebSocketController;
import com.opensource.docgrid.domain.rag.entity.RagResponse;
import com.opensource.docgrid.domain.rag.repository.RagResponseRepository;
-import com.opensource.docgrid.domain.search.enums.ResultStatus;
+import com.opensource.docgrid.domain.rag.service.command.RagResponseClaimService;
-import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
/**
- * PROCESSING 상태인 RagResponse를 하나씩 순서대로 꺼내 처리하는 경량 Worker (#218).
+ * PROCESSING 상태인 RagResponse를 최대 {@code rag.worker.max-concurrency}건까지 동시에 꺼내
+ * 처리하는 경량 Worker (#218, 병렬화는 #340).
*
- * {@code embedding_jobs}용 Worker(heartbeat·lease 복구 등 분산 처리 안전장치 포함, 26개 파일
- * 규모)와 달리, 이 Worker는 백엔드 인스턴스가 1개뿐이고 Ollama도 GPU 1개라 애초에 동시 처리가
- * 불가능하다는 전제 위에서 만들어졌다 — {@code @Scheduled} 폴링 하나로 충분하고, 여러 인스턴스
- * 간 조율(락·lease)은 필요 없다. Worker가 정확히 1개뿐이라는 사실 자체가 Ollama 호출의
- * 동시성 상한을 자연히 1로 만든다 — 기각했던 세마포어 게이트(#218 초안)가 하던 역할을 이
- * 구조가 대신한다.
+ * {@code embedding_jobs}용 Worker(heartbeat·lease 복구 등 분산 처리 안전장치 포함, 수십 파일
+ * 규모)와 달리, 이 Worker는 백엔드 인스턴스가 1개뿐이라는 전제 위에서 만들어졌다 — 여러 인스턴스
+ * 간 조율(락·lease)은 필요 없다. 다만 GPU/Ollama 하나가 실제로 감당 가능한 병렬 슬롯 수만큼은
+ * 이 프로세스 안에서 동시에 처리할 수 있다는 것이 #340의 전제다.
*
- * {@code processJob()} 실행(=OllamaClient HTTP 호출, 최대 {@code ollama.generate-deadline})이
- * 끝나야 다음 폴링이 실행되므로, 폴링 주기 자체는 혼잡 여부와 무관하게 큐가 밀리지 않는 한
- * 크게 중요하지 않다 — PROCESSING 건이 있으면 그 즉시 다음 턴에 잡힌다.
+ * 동시성은 두 계층으로 제한된다: ① {@link #ragWorkerSlots}(로컬 {@link Semaphore})를 먼저
+ * 확보해야 DB claim을 시도하고 — 이 순서 덕분에 "claim은 됐는데 실행할 스레드가 없는" 상태가
+ * 생기지 않는다. ② claim 자체는 {@link RagResponseClaimService}가 {@code FOR UPDATE SKIP
+ * LOCKED} + {@code claimed_at}으로 여러 스레드가 동시에 같은 job을 집지 못하게 막는다. 실제
+ * 처리는 {@link #ragWorkerJobExecutor}(전용 {@link ThreadPoolExecutor})에서 실행되어,
+ * {@link #processNext()} 자체는 claim만 하고 즉시 반환한다 — Ollama 호출(최대
+ * {@code ollama.generate-deadline})로 폴링 스레드가 막히지 않는다.
+ *
+ * {@link #recoverStaleClaimsOnStartup()}은 재시작 전 프로세스가 claim한 채 남긴 job의
+ * claim을 앱 시작 시 1회 풀어준다 — "인스턴스 1개" 전제를 유지하는 한 안전한 최소한의 복구다.
*/
@Component
-@RequiredArgsConstructor
@Slf4j
public class RagJobWorker {
private final RagResponseRepository ragResponseRepository;
+ private final RagResponseClaimService ragResponseClaimService;
private final RagFacade ragFacade;
private final RagWebSocketController ragWebSocketController;
+ private final Semaphore ragWorkerSlots;
+ private final ThreadPoolExecutor ragWorkerJobExecutor;
+
+ public RagJobWorker(
+ RagResponseRepository ragResponseRepository,
+ RagResponseClaimService ragResponseClaimService,
+ RagFacade ragFacade,
+ RagWebSocketController ragWebSocketController,
+ @Qualifier(RagExecutionConfig.RAG_WORKER_SLOTS) Semaphore ragWorkerSlots,
+ @Qualifier(RagExecutionConfig.RAG_WORKER_JOB_EXECUTOR) ThreadPoolExecutor ragWorkerJobExecutor
+ ) {
+ this.ragResponseRepository = ragResponseRepository;
+ this.ragResponseClaimService = ragResponseClaimService;
+ this.ragFacade = ragFacade;
+ this.ragWebSocketController = ragWebSocketController;
+ this.ragWorkerSlots = ragWorkerSlots;
+ this.ragWorkerJobExecutor = ragWorkerJobExecutor;
+ }
+
+ /**
+ * 앱 준비 완료 시 1회, 이전 프로세스가 claim한 채 완료하지 못한 job의 claim을 전부 풀어
+ * 재시작 뒤에도 다시 시도될 수 있게 한다(#340 CodeRabbit 리뷰 반영). 이 복구가 없으면
+ * {@code claimed_at}이 남아있는 job은 {@link RagResponseClaimService#claimNext}가 영원히
+ * 다시 집어주지 않아, 실제로 한 번도 재시도되지 않고 {@code RagJobTimeoutSweeper}의
+ * fallback만 기다리게 된다(#218 이전 방식은 이런 job을 자동으로 재시도했으므로 이 복구가
+ * 없으면 퇴보다). "인스턴스는 항상 1개"라는 전제 위에서만 안전 — 이 시점엔 다른 프로세스가
+ * 진짜로 처리 중일 수 없다.
+ */
+ @EventListener(ApplicationReadyEvent.class)
+ public void recoverStaleClaimsOnStartup() {
+ int recovered = ragResponseClaimService.recoverStaleClaimsOnStartup();
+ if (recovered > 0) {
+ log.warn("[RAG-WORKER] 재시작 복구: 이전 프로세스가 claim한 채 방치된 job {}건의 claim을 해제함", recovered);
+ }
+ }
+
+ /**
+ * 1초마다 실행되어, 로컬 슬롯이 남아있는 한 PROCESSING job을 계속 claim해 전용 Executor에
+ * 넘긴다. 슬롯이 없거나(이미 정원만큼 처리 중) 대기 중인 job이 없으면 그 자리에서 멈춘다.
+ *
+ * {@code while (tryAcquire())}만으로 반복 횟수 상한이 자동으로 걸린다 — Semaphore의 총
+ * permit 수가 이미 {@code max-concurrency}와 같아서, 별도 카운터 변수 없이도 이 루프가
+ * {@code max-concurrency}번보다 더 돌 수 없다.
+ */
+ @Scheduled(fixedDelayString = "${rag.worker.polling-interval:1s}")
+ public void processNext() {
+ while (ragWorkerSlots.tryAcquire()) {
+ Optional 처리 결과는 세 갈래로 갈린다: ①정상 성공 — WebSocket 알림. ②{@link
* OptimisticLockingFailureException} — 다른 트랜잭션이 이미 이 job을 처리했다는 뜻이라
@@ -51,49 +136,36 @@ public class RagJobWorker {
* boolean을 확인한 뒤에만 보낸다 — RagJobTimeoutSweeper가 이 job을 이미 먼저 FAILED로
* 확정해뒀다면(#288) 두 메서드 다 실제로는 아무것도 안 바꾸고 false를 반환하는데, 이 경우
* 스위퍼가 이미 보낸 알림 외에 Worker가 중복으로 또 보낼 이유가 없다.
+ *
+ * 어떤 경로로 끝나든 {@code finally}에서 반드시 슬롯을 반환한다 — 안 그러면 이 Worker가
+ * 처리 가능한 동시성이 영구히 줄어든다.
*/
- @Scheduled(fixedDelayString = "${rag.worker.polling-interval:1s}")
- public void processNext() {
- Optional {@code embedding_jobs}의 {@code EmbeddingJobClaimService}와 같은 트랜잭션 경계 전략을
+ * 쓴다 — 행 잠금({@code FOR UPDATE SKIP LOCKED})과 claim 표시를 하나의 짧은 트랜잭션으로
+ * 묶어 커밋과 동시에 락을 풀고, 실제 Ollama 호출은 이 트랜잭션 밖에서(별도 스레드가
+ * {@link com.opensource.docgrid.domain.rag.service.RagFacade#processJob}을 부르며) 진행한다.
+ * RAG는 PENDING 같은 별도 대기 상태가 없어(생성 즉시 PROCESSING) {@code embedding_jobs}처럼
+ * status 전이로 claim을 표시할 수 없다 — 대신 {@code claimed_at} 컬럼을 그 신호로 쓴다.
+ *
+ * {@code embedding_jobs}와 달리 별도의 분산 Worker 등록/생존 검증은 하지 않는다 — RAG
+ * Worker는 이 프로세스 안의 로컬 스레드일 뿐이라 그런 개념 자체가 없다.
+ */
+@Service
+@RequiredArgsConstructor
+@Transactional
+public class RagResponseClaimService {
+
+ private final RagResponseRepository ragResponseRepository;
+ private final Clock clock;
+
+ /**
+ * 다음으로 처리할 PROCESSING job 하나를 claim한다.
+ *
+ * @return claim에 성공한 job의 id. 대기 중인 job이 없으면 빈 값.
+ */
+ public Optional 병렬화(#340) 이후에는 여기에 "실제로 동시에 처리됐는가"까지 증명한다 — {@code claimed_at}이
+ * "언제부터 실제로 처리되기 시작했는지"를 알려주는 유일한 신호다({@code updatedAt}은 완료 확정이
+ * 전부 벌크 UPDATE라 채워지지 않는다). {@code rag.worker.max-concurrency}가 2 이상이면 3건 중
+ * 최소 2건은 거의 동시에 claim되어야 한다 — 셋 다 동시에 claim되길 요구하지 않는 이유는
+ * max-concurrency가 정확히 2일 때는 3번째 job이 앞선 두 건 중 하나가 끝날 때까지 자연스럽게
+ * 기다리기 때문이다(그래도 안전 실패는 없다 — #218의 핵심 목표).
*/
@Tag("integration")
@SpringBootTest
@@ -149,6 +156,19 @@ void threeConcurrentJobs_allEventuallySucceedViaRealScheduler() throws Interrupt
long successCount = finished.stream().filter(j -> j.getStatus() == ResultStatus.SUCCESS).count();
System.out.println("[TEST] SUCCESS=" + successCount + "/3, answers=" +
finished.stream().map(RagResponse::getAnswerText).toList());
+
+ // #340 병렬화 증명: claimed_at 3건 중 최소 2건은 서로 가까운 시각에 claim됐어야 한다 —
+ // 순차 처리였다면(#218 이전 방식) 각 claim은 앞선 job의 전체 처리 시간(수 초~수십 초)만큼
+ // 떨어져 있었을 것이다. max-concurrency가 정확히 2여도(기본값) 최소 두 건은 동시에 슬롯을
+ // 잡을 수 있으므로, "셋 다"가 아니라 "가장 가까운 두 건"의 간격으로 판단한다.
+ List 병렬화(#340) 이후 {@code processNext()}는 claim만 하고 실제 처리는 전용 Executor
+ * 스레드에 넘긴 뒤 즉시 반환한다 — 그래서 이 테스트도 호출 직후 동기적으로 결과를 확인하는
+ * 대신, {@link RagJobWorkerConcurrentQueueIntegrationTest}가 이미 쓰는 Awaitility로 처리가
+ * 끝날 때까지 기다린다.
*/
@Tag("integration")
@SpringBootTest
@@ -118,11 +125,15 @@ void processNext_persistsStatusChangeAcrossDetachedEntityBoundary() {
ragJobWorker.processNext();
- // Ollama가 로컬에 떠 있지 않을 수도 있으므로 SUCCESS/FAILED 둘 다 통과 조건으로 둔다 —
- // 이 테스트가 검증하는 건 "LLM 호출 성공 여부"가 아니라 "detached 상태에서도 최종
- // 상태가 DB에 반영되는지"다.
- RagResponse persisted = ragResponseRepository.findById(jobId).orElseThrow();
- assertThat(persisted.getStatus()).isNotEqualTo(ResultStatus.PROCESSING);
- assertThat(persisted.getAnswerText()).isNotNull();
+ // processNext()는 claim만 하고 즉시 반환하므로(#340), 실제 Ollama 호출·완료 확정은
+ // 전용 Executor 스레드에서 비동기로 이어진다 — 120초(read-timeout 90s + 여유)까지
+ // 기다렸다가 확인한다. Ollama가 로컬에 떠 있지 않을 수도 있으므로 SUCCESS/FAILED 둘 다
+ // 통과 조건으로 둔다 — 이 테스트가 검증하는 건 "LLM 호출 성공 여부"가 아니라 "detached
+ // 상태에서도 최종 상태가 DB에 반영되는지"다.
+ await().atMost(Duration.ofSeconds(120)).untilAsserted(() -> {
+ RagResponse persisted = ragResponseRepository.findById(jobId).orElseThrow();
+ assertThat(persisted.getStatus()).isNotEqualTo(ResultStatus.PROCESSING);
+ assertThat(persisted.getAnswerText()).isNotNull();
+ });
}
}
diff --git a/backend/src/test/java/com/opensource/docgrid/domain/rag/integration/RagResponseClaimIntegrationTest.java b/backend/src/test/java/com/opensource/docgrid/domain/rag/integration/RagResponseClaimIntegrationTest.java
new file mode 100644
index 00000000..73c0e25a
--- /dev/null
+++ b/backend/src/test/java/com/opensource/docgrid/domain/rag/integration/RagResponseClaimIntegrationTest.java
@@ -0,0 +1,251 @@
+package com.opensource.docgrid.domain.rag.integration;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import java.time.LocalDateTime;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.BrokenBarrierException;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.Supplier;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.TestInstance;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.test.context.SpringBootTest;
+import org.springframework.test.context.ActiveProfiles;
+import org.springframework.transaction.PlatformTransactionManager;
+import org.springframework.transaction.TransactionDefinition;
+import org.springframework.transaction.support.TransactionTemplate;
+
+import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel;
+import com.opensource.docgrid.domain.embedding.fixture.EmbeddingModelFixture;
+import com.opensource.docgrid.domain.embedding.repository.EmbeddingModelRepository;
+import com.opensource.docgrid.domain.rag.entity.RagResponse;
+import com.opensource.docgrid.domain.rag.repository.RagResponseRepository;
+import com.opensource.docgrid.domain.rag.service.command.RagResponseClaimService;
+import com.opensource.docgrid.domain.search.entity.SearchConversation;
+import com.opensource.docgrid.domain.search.entity.SearchQuery;
+import com.opensource.docgrid.domain.search.enums.ResultStatus;
+import com.opensource.docgrid.domain.search.enums.SearchType;
+import com.opensource.docgrid.domain.search.repository.SearchConversationRepository;
+import com.opensource.docgrid.domain.search.repository.SearchQueryRepository;
+import com.opensource.docgrid.domain.user.entity.User;
+import com.opensource.docgrid.domain.user.enums.UserStatus;
+import com.opensource.docgrid.domain.user.repository.UserRepository;
+
+/**
+ * 실제 PostgreSQL에서 RAG Job Claim의 행 잠금·동시 claim 불변식을 검증하는 통합 테스트 (#340).
+ *
+ * {@code EmbeddingJobClaimIntegrationTest}와 같은 방식으로, 서로 다른 Thread와
+ * {@code REQUIRES_NEW} Transaction을 사용해 단일 Persistence Context의 순차 호출로는 재현할 수
+ * 없는 {@code FOR UPDATE SKIP LOCKED} 경쟁을 검증한다.
+ */
+@Tag("integration")
+@SpringBootTest
+@ActiveProfiles("test")
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+@DisplayName("RagResponse Claim DB 동시성 통합 테스트")
+class RagResponseClaimIntegrationTest {
+
+ private static final long TIMEOUT_SECONDS = 10;
+
+ @Autowired private PlatformTransactionManager transactionManager;
+ @Autowired private RagResponseRepository ragResponseRepository;
+ @Autowired private RagResponseClaimService ragResponseClaimService;
+ @Autowired private SearchQueryRepository searchQueryRepository;
+ @Autowired private SearchConversationRepository searchConversationRepository;
+ @Autowired private UserRepository userRepository;
+ @Autowired private EmbeddingModelRepository embeddingModelRepository;
+
+ private ExecutorService executorService;
+ private final AtomicInteger threadSequence = new AtomicInteger();
+ private final List {@code Semaphore}는 mock하지 않고 실제 인스턴스를 쓴다 — I/O가 없는 순수 카운터라, mock보다
+ * 실제 객체로 "슬롯이 진짜 반환됐는지"를 permit 개수로 직접 확인하는 편이 더 간단하고 정확하다.
*/
@ExtendWith(MockitoExtension.class)
@DisplayName("RagJobWorker 단위 테스트")
class RagJobWorkerTest {
- @InjectMocks
- private RagJobWorker ragJobWorker;
-
@Mock
private RagResponseRepository ragResponseRepository;
+ @Mock
+ private RagResponseClaimService ragResponseClaimService;
+
@Mock
private RagFacade ragFacade;
@Mock
private RagWebSocketController ragWebSocketController;
+ @Mock
+ private ThreadPoolExecutor ragWorkerJobExecutor;
+
+ private Semaphore ragWorkerSlots;
+ private RagJobWorker ragJobWorker;
+
+ @BeforeEach
+ void setUp() {
+ ragWorkerSlots = new Semaphore(1);
+ ragJobWorker = new RagJobWorker(
+ ragResponseRepository, ragResponseClaimService, ragFacade, ragWebSocketController,
+ ragWorkerSlots, ragWorkerJobExecutor
+ );
+ }
+
@Test
- @DisplayName("PROCESSING 건이 없으면 아무것도 하지 않는다")
- void processNext_noPendingJob_doesNothing() {
- given(ragResponseRepository.findFirstByStatusOrderByCreatedAtAsc(ResultStatus.PROCESSING))
- .willReturn(Optional.empty());
+ @DisplayName("앱 시작 시 복구할 claim이 있으면 경고 로그를 남긴다(부작용은 claim 서비스에 위임)")
+ void recoverStaleClaimsOnStartup_delegatesToClaimServiceAndReturns() {
+ given(ragResponseClaimService.recoverStaleClaimsOnStartup()).willReturn(2);
+
+ ragJobWorker.recoverStaleClaimsOnStartup();
+
+ then(ragResponseClaimService).should(times(1)).recoverStaleClaimsOnStartup();
+ }
+
+ @Test
+ @DisplayName("슬롯이 없으면 claim 자체를 시도하지 않는다")
+ void processNext_noSlotAvailable_neverClaims() {
+ // 이미 다른 job이 유일한 슬롯을 쓰고 있는 상황을 재현한다. permit이 이미 있는 상태의
+ // acquireUninterruptibly()는 즉시 반환되므로 실제로 블로킹되지 않는다.
+ ragWorkerSlots.acquireUninterruptibly();
ragJobWorker.processNext();
- then(ragFacade).should(never()).processJob(any());
- then(ragWebSocketController).should(never()).notifyAnswerReady(any(), any());
+ then(ragResponseClaimService).should(never()).claimNext();
+ then(ragWorkerJobExecutor).should(never()).execute(any());
}
@Test
- @DisplayName("PROCESSING 건이 있으면 처리하고, 요청자 본인에게만 완료를 push한다")
- void processNext_pendingJobExists_processesAndNotifiesOwner() {
+ @DisplayName("claim할 job이 없으면 슬롯을 반환하고 Executor를 부르지 않는다")
+ void processNext_claimEmpty_releasesSlotAndSkipsExecutor() {
+ given(ragResponseClaimService.claimNext()).willReturn(Optional.empty());
+
+ ragJobWorker.processNext();
+
+ then(ragWorkerJobExecutor).should(never()).execute(any());
+ assertThat(ragWorkerSlots.availablePermits()).isEqualTo(1);
+ }
+
+ @Test
+ @DisplayName("claim에 성공하면 Executor에 제출하고, 정상 처리되면 요청자 본인에게만 완료를 push한다")
+ void processNext_claimSucceeds_submitsAndNotifiesOwner() {
RagResponse job = deepStubJob(999L, 100L, "user@example.com");
- given(ragResponseRepository.findFirstByStatusOrderByCreatedAtAsc(ResultStatus.PROCESSING))
- .willReturn(Optional.of(job));
+ given(ragResponseClaimService.claimNext()).willReturn(Optional.of(999L));
+ given(ragResponseRepository.findWithQueryAndUserById(999L)).willReturn(Optional.of(job));
given(ragFacade.processJob(999L)).willReturn(true);
ragJobWorker.processNext();
+ runSubmittedTask();
- // Worker는 detached entity를 그대로 넘기지 않고 id만 넘긴다 — processJob()이 자기 트랜잭션
- // 안에서 다시 조회해야 완료 처리(조건부 UPDATE)가 최신 상태 기준으로 실행된다.
then(ragFacade).should(times(1)).processJob(999L);
then(ragWebSocketController).should(times(1)).notifyAnswerReady("user@example.com", 100L);
+ // 처리(성공적으로 실행된 Runnable)가 끝나면 finally에서 슬롯을 되돌려준다.
+ assertThat(ragWorkerSlots.availablePermits()).isEqualTo(1);
}
@Test
@DisplayName("경합(#288): processJob이 false를 반환하면(RagJobTimeoutSweeper가 이미 확정함) 알림을 보내지 않는다")
void processNext_processJobLosesRace_doesNotNotify() {
RagResponse job = deepStubJob(999L, 100L, "user@example.com");
- given(ragResponseRepository.findFirstByStatusOrderByCreatedAtAsc(ResultStatus.PROCESSING))
- .willReturn(Optional.of(job));
+ given(ragResponseClaimService.claimNext()).willReturn(Optional.of(999L));
+ given(ragResponseRepository.findWithQueryAndUserById(999L)).willReturn(Optional.of(job));
given(ragFacade.processJob(999L)).willReturn(false);
ragJobWorker.processNext();
+ runSubmittedTask();
then(ragWebSocketController).should(never()).notifyAnswerReady(any(), any());
}
@Test
- @DisplayName("processJob이 예상 밖 예외를 던지면 job을 FAILED로 확정하고, Worker는 죽지 않고 이번 건만 건너뛴다")
- void processNext_unexpectedException_marksFailedAndSkipsJobWithoutCrashingWorker() {
+ @DisplayName("processJob이 예상 밖 예외를 던지면 job을 FAILED로 확정하고, 알림은 그대로 보낸다")
+ void processNext_unexpectedException_marksFailedAndNotifies() {
RagResponse job = deepStubJob(999L, 100L, "user@example.com");
- given(ragResponseRepository.findFirstByStatusOrderByCreatedAtAsc(ResultStatus.PROCESSING))
- .willReturn(Optional.of(job));
- org.mockito.Mockito.doThrow(new RuntimeException("예상 밖 버그")).when(ragFacade).processJob(999L);
+ given(ragResponseClaimService.claimNext()).willReturn(Optional.of(999L));
+ given(ragResponseRepository.findWithQueryAndUserById(999L)).willReturn(Optional.of(job));
+ doThrow(new RuntimeException("예상 밖 버그")).when(ragFacade).processJob(999L);
given(ragFacade.markUnexpectedFailure(999L, "예상 밖 버그")).willReturn(true);
ragJobWorker.processNext();
+ runSubmittedTask();
- // job을 PROCESSING으로 방치하면 Worker가 같은 job을 계속 다시 집어 무한 재시도하게 된다
- // (detached entity 버그와 같은 증상) — 그래서 반드시 FAILED로 확정해야 한다.
then(ragFacade).should(times(1)).markUnexpectedFailure(999L, "예상 밖 버그");
- // FAILED로 확정된 이상 사용자도 결과(비록 실패 안내지만)를 받아야 하므로 알림은 그대로 간다.
then(ragWebSocketController).should(times(1)).notifyAnswerReady("user@example.com", 100L);
}
@@ -109,44 +157,61 @@ void processNext_unexpectedException_marksFailedAndSkipsJobWithoutCrashingWorker
@DisplayName("다른 트랜잭션이 이미 같은 job을 처리했으면(낙관적 락 경합) FAILED로 덮어쓰지 않고 조용히 넘어간다")
void processNext_optimisticLockingFailure_skipsWithoutOverwritingAsFailed() {
RagResponse job = deepStubJob(999L, 100L, "user@example.com");
- given(ragResponseRepository.findFirstByStatusOrderByCreatedAtAsc(ResultStatus.PROCESSING))
- .willReturn(Optional.of(job));
- org.mockito.Mockito.doThrow(new OptimisticLockingFailureException("경합"))
- .when(ragFacade).processJob(999L);
+ given(ragResponseClaimService.claimNext()).willReturn(Optional.of(999L));
+ given(ragResponseRepository.findWithQueryAndUserById(999L)).willReturn(Optional.of(job));
+ doThrow(new OptimisticLockingFailureException("경합")).when(ragFacade).processJob(999L);
ragJobWorker.processNext();
+ runSubmittedTask();
- // 다른 트랜잭션이 이미 올바르게 처리한 결과이므로, 이걸 FAILED로 덮어쓰면 정상 처리된
- // 결과를 오답으로 바꿔버리는 2차 사고가 난다 — markUnexpectedFailure를 호출하면 안 된다.
then(ragFacade).should(never()).markUnexpectedFailure(any(), any());
then(ragWebSocketController).should(never()).notifyAnswerReady(any(), any());
}
@Test
- @DisplayName("한 job이 예외로 실패해도 다음 폴링에서 뒤에 대기 중인 job이 정상 처리된다")
- void processNext_afterUnexpectedFailure_nextPollingProcessesFollowingJob() {
- RagResponse failingJob = deepStubJob(1L, 100L, "user1@example.com");
- RagResponse nextJob = deepStubJob(2L, 200L, "user2@example.com");
- org.mockito.Mockito.doThrow(new RuntimeException("예상 밖 버그")).when(ragFacade).processJob(1L);
-
- given(ragResponseRepository.findFirstByStatusOrderByCreatedAtAsc(ResultStatus.PROCESSING))
- .willReturn(Optional.of(failingJob));
- ragJobWorker.processNext(); // 1번째 폴링: failingJob 실패 → FAILED로 확정됨
-
- // FAILED로 확정됐으니 실제 DB에선 이제 findFirst...가 다음 대기 건(nextJob)을 돌려준다 —
- // 여기서는 그 상태 변화를 목으로 흉내낸다.
- given(ragResponseRepository.findFirstByStatusOrderByCreatedAtAsc(ResultStatus.PROCESSING))
- .willReturn(Optional.of(nextJob));
- given(ragFacade.processJob(2L)).willReturn(true);
- ragJobWorker.processNext(); // 2번째 폴링: nextJob은 정상 처리돼야 한다
-
- then(ragFacade).should(times(1)).processJob(2L);
- then(ragWebSocketController).should(times(1)).notifyAnswerReady("user2@example.com", 200L);
+ @DisplayName("claim 중 예외가 나면 슬롯을 반환하고 Executor를 부르지 않는다")
+ void processNext_claimThrows_releasesSlotAndSkipsExecutor() {
+ given(ragResponseClaimService.claimNext()).willThrow(new RuntimeException("DB 오류"));
+
+ ragJobWorker.processNext();
+
+ then(ragWorkerJobExecutor).should(never()).execute(any());
+ assertThat(ragWorkerSlots.availablePermits()).isEqualTo(1);
+ }
+
+ @Test
+ @DisplayName("Executor 제출이 거부되면(RejectedExecutionException) 슬롯을 반환한다")
+ void processNext_executorRejects_releasesSlot() {
+ given(ragResponseClaimService.claimNext()).willReturn(Optional.of(999L));
+ doThrow(new RejectedExecutionException()).when(ragWorkerJobExecutor).execute(any());
+
+ ragJobWorker.processNext();
+
+ assertThat(ragWorkerSlots.availablePermits()).isEqualTo(1);
+ }
+
+ @Test
+ @DisplayName("claim 직후 job이 사라졌으면(극단적 상황) processJob을 부르지 않고 슬롯만 반환한다")
+ void processNext_claimedJobVanished_skipsProcessing() {
+ given(ragResponseClaimService.claimNext()).willReturn(Optional.of(999L));
+ given(ragResponseRepository.findWithQueryAndUserById(999L)).willReturn(Optional.empty());
+
+ ragJobWorker.processNext();
+ runSubmittedTask();
+
+ then(ragFacade).should(never()).processJob(any());
+ assertThat(ragWorkerSlots.availablePermits()).isEqualTo(1);
+ }
+
+ /** Executor에 제출된 Runnable을 캡처해 테스트 스레드에서 즉시(동기) 실행한다. */
+ private void runSubmittedTask() {
+ ArgumentCaptor