Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
bae9acd
feat: #218 비동기 RAG 처리를 위한 데이터 계층 준비
kangcheolung Aug 17, 2026
ef68ba5
feat: #218 RagFacade를 enqueue/processJob으로 분리
kangcheolung Aug 17, 2026
2e8b42d
feat: #218 RAG Job Worker 추가
kangcheolung Aug 17, 2026
2cb6a0b
feat: #218 RAG 답변 완료 WebSocket 알림
kangcheolung Aug 17, 2026
04300b7
feat: #218 검색-RAG API를 비동기 계약으로 전환
kangcheolung Aug 17, 2026
d709863
chore: #218 Ollama 타임아웃 설정을 비동기 전환에 맞게 재조정
kangcheolung Aug 17, 2026
806f696
test: #218 비동기 RAG 처리 단위·통합 테스트
kangcheolung Aug 17, 2026
b0d9bd6
feat: #218 검색 결과 먼저 표시하고 AI 답변은 비동기로 갱신하는 UI
kangcheolung Aug 17, 2026
6f5e2ac
test: #218 SearchResponse ragStatus 필드 추가에 맞춰 프론트 테스트 갱신
kangcheolung Aug 17, 2026
9ff4a62
docs: #218 비동기 RAG Job 큐 설계 문서
kangcheolung Aug 17, 2026
aefee9b
docs: #35, #44 설계 문서에 이후 변경 이력 추가
kangcheolung Aug 17, 2026
4b939d5
fix: #218 CodeRabbit 지적 반영 — Worker 경합·예외 처리, 프론트 stale response 방지
kangcheolung Aug 17, 2026
7333d71
test: #218 CodeRabbit 지적 반영 — Worker 예외/경합 처리 회귀 테스트
kangcheolung Aug 17, 2026
b0261c7
docs: #218 CodeRabbit 지적 반영 — 코드 펜스 언어 태그 지정
kangcheolung Aug 17, 2026
ee28d82
docs: #218 CodeRabbit 리뷰로 고친 버그 3건 설계 문서에 반영
kangcheolung Aug 17, 2026
File filter

Filter by extension

Filter by extension

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

import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.EnableScheduling;

/**
* {@code WorkerSchedulingConfig} 등 다른 도메인의 {@code @EnableScheduling}과 별개로 켠다 —
* 그쪽은 {@code indexing.worker.enabled} 조건부라 꺼질 수 있지만, {@link
* com.opensource.docgrid.domain.rag.service.RagJobWorker}는 검색 API의 핵심 경로라 조건 없이
* 항상 돌아야 한다. {@code @EnableScheduling}을 여러 설정 클래스에 중복 선언해도 Spring이
* 안전하게 병합하므로 문제없다.
*/
@Configuration
@EnableScheduling
public class RagSchedulingConfig {
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
package com.opensource.docgrid.domain.rag.controller;

import org.springframework.messaging.simp.SimpMessagingTemplate;
import org.springframework.stereotype.Component;

import lombok.RequiredArgsConstructor;

/**
* RAG 답변이 준비됐음을 요청한 사용자에게만 push하는 전송 계층 (#218).
*
* <p>{@code DashboardWebSocketController}(/topic/dashboard, 전체 브로드캐스트)와 달리, 이건 검색
* 요청을 보낸 그 유저 한 명에게만 전달돼야 한다 — 관리자 전용 브로드캐스트 채널을 재사용할 수 없는
* 이유다. {@code convertAndSendToUser}는 {@code StompAuthChannelInterceptor}가 CONNECT 시점에
* 세션에 부착한 Principal(이메일)로 목적지를 사용자별로 격리한다 — 다른 유저는 같은 목적지
* ({@code /user/queue/rag-answer})를 구독해도 이 메시지를 받지 않으므로, 대시보드처럼 별도의
* 구독 인가 Interceptor가 필요 없다.
*
* <p>본문은 트리거 용도로만 쓴다. {@code useDashboardSocket}과 동일하게, 프론트는 이 메시지를
* "다시 조회해야 한다"는 신호로만 쓰고 최신 상태는 REST로 다시 읽는다 — Push 페이로드와 실제
* DB 상태가 어긋날 걱정 없이 항상 단일 진실 소스(REST)를 신뢰할 수 있다.
*/
@Component
@RequiredArgsConstructor
public class RagWebSocketController {

private static final String RAG_ANSWER_QUEUE = "/queue/rag-answer";

private final SimpMessagingTemplate messagingTemplate;

public void notifyAnswerReady(String userEmail, Long queryId) {
messagingTemplate.convertAndSendToUser(userEmail, RAG_ANSWER_QUEUE, new RagAnswerReadyEvent(queryId));
}

// 완료 알림의 최소 트리거 페이로드 — 답변 본문은 담지 않는다. 프론트가 이 이벤트를 받으면
// 항상 GET /search/{queryId}로 다시 조회해야 하며, 이 record 자체를 최종 상태로 신뢰하면 안 된다.
private record RagAnswerReadyEvent(Long queryId) {
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
package com.opensource.docgrid.domain.rag.dto;

/**
* RagFacade.enqueue() 반환 타입 — 검색 후보가 없어(NO_CONTEXT) LLM 호출 없이 즉시 끝난 경우와,
* PROCESSING으로 Job 큐에 올라가 Worker의 처리를 기다려야 하는 경우를 구분한다.
*/
public record RagEnqueueOutcome(RagAnswer immediateAnswer, boolean pending) {

public static RagEnqueueOutcome stillPending() {
return new RagEnqueueOutcome(null, true);
}

public static RagEnqueueOutcome done(RagAnswer answer) {
return new RagEnqueueOutcome(answer, false);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,8 @@ public class RagResponse extends BaseEntity {
@JoinColumn(name = "query_id", nullable = false)
private SearchQuery query;

@Column(name = "answer_text", nullable = false, columnDefinition = "TEXT")
// PROCESSING 상태로 처음 저장될 때는 아직 값이 없다 — Worker가 생성을 마치면 markSuccess/markFailed로 채운다.
@Column(name = "answer_text", columnDefinition = "TEXT")
private String answerText;

// LLM 제공자 이름
Expand Down Expand Up @@ -99,4 +100,23 @@ public RagResponse(SearchQuery query, String answerText, String llmProvider, Str
this.status = status;
this.errorMessage = errorMessage;
}

// Worker가 LLM 생성을 마친 뒤 PROCESSING 상태였던 이 row를 SUCCESS로 채운다.
public void markSuccess(String answerText, String llmModelName, Integer inputTokenCount,
Integer outputTokenCount, Integer latencyMs) {
this.answerText = answerText;
this.llmModelName = llmModelName;
this.inputTokenCount = inputTokenCount;
this.outputTokenCount = outputTokenCount;
this.latencyMs = latencyMs;
this.status = ResultStatus.SUCCESS;
}

// LLM 호출 실패 시에도 빈손이 아니라 extractive fallback 답변을 채워 넣는다 — status만 FAILED로
// 남겨 감사 추적을 위한 실패 이력은 유지한다.
public void markFailed(String fallbackAnswerText, String errorMessage) {
this.answerText = fallbackAnswerText;
this.status = ResultStatus.FAILED;
this.errorMessage = errorMessage;
}
}
Original file line number Diff line number Diff line change
@@ -1,8 +1,20 @@
package com.opensource.docgrid.domain.rag.repository;

import java.util.Optional;

import org.springframework.data.jpa.repository.EntityGraph;
import org.springframework.data.jpa.repository.JpaRepository;

import com.opensource.docgrid.domain.rag.entity.RagResponse;
import com.opensource.docgrid.domain.search.enums.ResultStatus;

public interface RagResponseRepository extends JpaRepository<RagResponse, Long> {

// Worker가 순서대로 하나씩 꺼내 처리한다 — Worker가 1개뿐이라 별도 락/claim 없이도 안전하다.
// query/query.user를 미리 fetch해 Worker가 트랜잭션 밖(WebSocket push 시점)에서
// job.getQuery().getUser().getEmail()에 접근해도 LazyInitializationException이 나지 않게 한다.
@EntityGraph(attributePaths = {"query", "query.user"})
Optional<RagResponse> findFirstByStatusOrderByCreatedAtAsc(ResultStatus status);
Comment thread
coderabbitai[bot] marked this conversation as resolved.

Optional<RagResponse> findByQuery_Id(Long queryId);
}
Original file line number Diff line number Diff line change
@@ -1,8 +1,12 @@
package com.opensource.docgrid.domain.rag.repository;

import java.util.List;

import org.springframework.data.jpa.repository.JpaRepository;

import com.opensource.docgrid.domain.rag.entity.ResponseCitation;

public interface ResponseCitationRepository extends JpaRepository<ResponseCitation, Long> {

List<ResponseCitation> findByResponse_IdOrderByCitationOrder(Long responseId);
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,13 +7,17 @@

import com.opensource.docgrid.domain.rag.dto.OllamaGenerateResult;
import com.opensource.docgrid.domain.rag.dto.RagAnswer;
import com.opensource.docgrid.domain.rag.dto.RagEnqueueOutcome;
import com.opensource.docgrid.domain.rag.entity.RagResponse;
import com.opensource.docgrid.domain.rag.repository.RagResponseRepository;
import com.opensource.docgrid.domain.rag.service.command.RagResponseCommandService;
import com.opensource.docgrid.domain.rag.service.command.ResponseCitationCommandService;
import com.opensource.docgrid.domain.search.dto.VectorSearchCandidate;
import com.opensource.docgrid.domain.search.entity.SearchQuery;
import com.opensource.docgrid.domain.search.entity.SearchResult;
import com.opensource.docgrid.domain.search.repository.SearchResultRepository;
import com.opensource.docgrid.global.exception.DocGridException;
import com.opensource.docgrid.global.exception.ErrorCode;

import jakarta.persistence.EntityManager;
import lombok.RequiredArgsConstructor;
Expand All @@ -22,16 +26,16 @@
/**
* RAG 답변 생성 전체 흐름을 조율하는 Facade (F-RAG-05).
*
* <p>비동기 Job 큐 전환(#218) 이후 두 단계로 나뉜다:
* <pre>
* 1. candidates가 비어있으면(NO_CONTEXT) LLM 호출 없이 고정 응답 저장
* 2. PromptBuilder로 프롬프트 조립
* 3. OllamaClient 호출
* 4. rag_responses 저장 (성공/실패)
* 5. 성공 시 response_citations 저장, 실패 시 검색 후보와 안내 답변 반환
* 1. enqueue() — SearchController가 검색 직후 동기 호출. 프롬프트만 조립해 PROCESSING으로 저장하고
* 즉시 반환한다(LLM 호출 없음). candidates가 비어있으면(NO_CONTEXT) 여기서 바로 끝난다.
* 2. processJob() — RagJobWorker가 PROCESSING row를 하나씩 꺼내 호출. 실제 OllamaClient 호출과
* 결과 영속화(rag_responses, response_citations)를 담당한다.
* </pre>
*
* <p>SearchFacade와 별도 트랜잭션으로 분리되어 있다(SearchController가 순차 호출) — 검색 DB 작업과
* LLM HTTP 호출을 포함한 RAG DB 작업이 하나의 커넥션을 오래 물고 있지 않도록 하기 위함이다.
* <p>SearchFacade와 별도 트랜잭션으로 분리되어 있다(SearchController가 순차 호출) — 검색 DB 작업이
* enqueue()의 짧은 DB 작업과 하나의 커넥션을 오래 물고 있지 않도록 하기 위함이다.
*/
@Transactional
@Service
Expand All @@ -42,12 +46,17 @@ public class RagFacade {
private static final String LLM_FALLBACK_PREFIX = "AI 답변 생성이 지연되고 있습니다. "
+ "가장 관련도 높은 문서에서 다음 내용을 찾았습니다:\n\n";

// processJob() 내부에서 예상 못한 예외(버그 등)로 실패했을 때 쓰는 최소 안내 문구. extractive
// fallback과 달리 candidates를 다시 불러오지 않는다 — 이미 한 번 예상 밖으로 실패한 상황에서
// 추가 조회를 시도하다 또 실패할 위험을 만들지 않기 위함이다(RagJobWorker 참고).
private static final String UNEXPECTED_FAILURE_ANSWER_TEXT = "답변 생성 중 예상치 못한 오류가 발생했습니다.";

// fallback 문구에 원문을 통째로 붙이면 답변이 지나치게 길어져, 미리보기 수준으로만 잘라 보여준다.
private static final int FALLBACK_EXCERPT_MAX_CODE_POINTS = 300;

// topK는 호출자가 1~20까지 자유롭게 요청할 수 있어(SearchRequest), 후보 수를 그대로 프롬프트에
// 다 넣으면 prefill 시간이 예측 불가능해져 read-timeout(25s)을 넘기는 경우가 생긴다.
// 화면에 보여줄 인용 문서 수(topK)와 별개로, LLM이 실제로 읽는 후보 수는 이 값으로 고정한다.
// 다 넣으면 prefill 시간이 예측 불가능해진다. 화면에 보여줄 인용 문서 수(topK)와 별개로,
// LLM이 실제로 읽는 후보 수는 이 값으로 고정한다.
private static final int MAX_PROMPT_CANDIDATES = 3;

// PromptBuilder가 LLM에게 무관한 문서일 때 이 문구로만 답하도록 지시한다 — 검색은 됐지만(candidates
Expand All @@ -58,56 +67,99 @@ public class RagFacade {
private final OllamaClient ollamaClient;
private final RagResponseCommandService ragResponseCommandService;
private final ResponseCitationCommandService responseCitationCommandService;
private final RagResponseRepository ragResponseRepository;
private final SearchResultRepository searchResultRepository;
private final EntityManager entityManager;

public RagAnswer generate(
Long queryId, String queryText, List<VectorSearchCandidate> candidates, List<SearchResult> searchResults
) {
public RagEnqueueOutcome enqueue(Long queryId, String queryText, List<VectorSearchCandidate> candidates) {
SearchQuery queryRef = entityManager.getReference(SearchQuery.class, queryId);

// 검색 후보가 없으면(NO_CONTEXT) LLM 호출 없이 고정 응답 저장
// 검색 후보가 없으면(NO_CONTEXT) LLM 호출 없이 고정 응답으로 바로 끝낸다 — Job 큐에 올릴 이유가 없다.
if (candidates.isEmpty()) {
RagResponse ragResponse = ragResponseCommandService.createNoContext(queryRef);
log.info("[RAG] no context queryId={} responseId={}", queryId, ragResponse.getId());
return RagAnswer.noContext(ragResponse.getAnswerText());
return RagEnqueueOutcome.done(RagAnswer.noContext(ragResponse.getAnswerText()));
}

// 검색 후보가 있으면 프롬프트 조립 후 LLM 호출 (LLM 입력은 상위 MAX_PROMPT_CANDIDATES개로 제한)
List<VectorSearchCandidate> promptCandidates = candidates.size() > MAX_PROMPT_CANDIDATES
? candidates.subList(0, MAX_PROMPT_CANDIDATES)
: candidates;
String prompt = promptBuilder.build(queryText, promptCandidates);
RagResponse ragResponse = ragResponseCommandService.createPending(queryRef, prompt);
log.info("[RAG] enqueued queryId={} responseId={}", queryId, ragResponse.getId());
return RagEnqueueOutcome.stillPending();
}

// RagJobWorker가 findFirstByStatusOrderByCreatedAtAsc()로 꺼낸 job은 그 조회 시점에 트랜잭션이
// 끝나 detached 상태다 — 그 인스턴스를 그대로 받아 markSuccess/markFailed로 값을 바꿔도 이
// 메서드의 새 트랜잭션에서는 dirty checking이 감지하지 못해 DB에 반영되지 않는다(영원히
// PROCESSING으로 남아 Worker가 같은 job을 계속 재처리하는 버그로 이어졌었다). 그래서 id만 받아
// 이 메서드 자신의 트랜잭션 안에서 다시 조회해 반드시 managed 상태로 확보한다.
public void processJob(Long jobId) {
RagResponse job = ragResponseRepository.findById(jobId)
.orElseThrow(() -> new DocGridException(ErrorCode.RAG_ANSWER_NOT_FOUND));
Long queryId = job.getQuery().getId();

OllamaGenerateResult result;
try {
result = ollamaClient.generate(prompt);
result = ollamaClient.generate(job.getPromptText());
} catch (DocGridException e) {
ragResponseCommandService.createFailed(queryRef, prompt, e.getMessage());
// LLM 장애가 권한 검증을 통과한 벡터 검색 결과까지 숨기지 않도록, 최상위 후보 원문을
// 그대로 인용해 최소한의 답을 제공한다(extractive fallback).
// 그대로 인용해 최소한의 답을 제공한다(extractive fallback). 이 fallback은 비동기 전환
// 이전과 달리 rag_responses에 그대로 영속화된다 — 나중에 GET/조회로 이 값을 그대로 돌려준다.
List<VectorSearchCandidate> candidates = loadCandidates(queryId);
String fallbackAnswer = candidates.isEmpty() ? e.getErrorCode().getMessage()
: buildExtractiveFallbackAnswer(candidates);
ragResponseCommandService.completeFailed(job, fallbackAnswer, e.getMessage());
log.warn("[RAG] fallback queryId={} errorCode={}", queryId, e.getErrorCode().getCode());
return RagAnswer.of(buildExtractiveFallbackAnswer(candidates), candidates);
return;
}

// LLM 이후의 영속화 실패는 검색 저하 응답으로 숨기지 않고 Transaction 오류로 전달한다.
RagResponse ragResponse = ragResponseCommandService.createSuccess(queryRef, prompt, result);
responseCitationCommandService.saveAll(ragResponse, candidates, searchResults);
log.info("[RAG] done queryId={} responseId={} latencyMs={}", queryId, ragResponse.getId(), result.latencyMs());

// LLM이 무관하다고 판단해 안내 문구로만 답했으면, 후보 문서를 근거처럼 같이 보여주지 않는다.
// 단, 7B 모델이 정상 답변을 끝낸 뒤 지시문을 메아리처럼 이 문구를 덧붙이는 패턴이 관찰됨 —
// 문구가 답변의 사실상 전부(맨 앞)일 때만 무관으로 취급하고, 정상 답변 중간에 박힌 문구는
// 그 지점부터 잘라내고 근거 문서는 유지한다.
// LLM이 무관하다고 판단해 안내 문구로만 답했으면, 근거 문서를 같이 보여주지 않는다. 단, 7B
// 모델이 정상 답변을 끝낸 뒤 지시문을 메아리처럼 이 문구를 덧붙이는 패턴이 관찰됨 — 문구가
// 답변의 사실상 전부(맨 앞)일 때만 무관으로 취급하고, 정상 답변 중간에 박힌 문구는 그
// 지점부터 잘라내고 근거 문서는 유지한다. 잘라낸 결과를 그대로 영속화해야 GET 조회 시
// 사용자에게 보이는 값과 DB 값이 일치한다(동기 시절엔 반환값에만 트리밍이 적용되고 DB엔
// 원문이 남았는데, 비동기에서는 이 row가 유일한 진실 소스라 그대로 두면 안 된다).
String answerText = result.answerText();
List<VectorSearchCandidate> candidates = loadCandidates(queryId);
boolean noRelevant = false;
int phraseIndex = answerText != null ? answerText.indexOf(NO_RELEVANT_DOC_PHRASE) : -1;
if (phraseIndex >= 0) {
if (answerText.strip().startsWith(NO_RELEVANT_DOC_PHRASE)) {
return RagAnswer.of(answerText, List.of());
noRelevant = true;
} else {
log.warn("[RAG] 정상 답변에 무관 안내 문구 혼입, 해당 지점부터 제거: queryId={} phraseIndex={}",
queryId, phraseIndex);
answerText = answerText.substring(0, phraseIndex).strip();
}
log.warn("[RAG] 정상 답변에 무관 안내 문구 혼입, 해당 지점부터 제거: queryId={} phraseIndex={}",
queryId, phraseIndex);
answerText = answerText.substring(0, phraseIndex).strip();
}
return RagAnswer.of(answerText, candidates);

ragResponseCommandService.completeSuccess(job, new OllamaGenerateResult(
result.model(), answerText, result.inputTokenCount(), result.outputTokenCount(), result.latencyMs()
));

if (!noRelevant) {
List<SearchResult> searchResults = searchResultRepository.findByQuery_IdOrderByRankNo(queryId);
responseCitationCommandService.saveAll(job, candidates, searchResults);
}
log.info("[RAG] done queryId={} responseId={} latencyMs={}", queryId, job.getId(), result.latencyMs());
}

// RagJobWorker가 processJob() 호출 중 예상 못한 예외(버그 등)를 잡았을 때 호출한다. 여기서
// FAILED로 확정하지 않으면 job이 영원히 PROCESSING으로 남아, 같은 job을 Worker가 계속
// 다시 집어 무한 재시도하게 된다 — 4-5에서 고친 detached entity 버그와 증상이 같아진다.
public void markUnexpectedFailure(Long jobId, String errorMessage) {
ragResponseRepository.findById(jobId)
.ifPresent(job -> ragResponseCommandService.completeFailed(job, UNEXPECTED_FAILURE_ANSWER_TEXT, errorMessage));
}

// Worker는 검색 시점의 in-memory candidates를 갖고 있지 않으므로, 이미 영속화된 search_results(+chunk)에서
// 동일한 순서로 다시 조립한다 — PromptBuilder에 넘겼던 것과 citation_order가 어긋나지 않는다.
private List<VectorSearchCandidate> loadCandidates(Long queryId) {
return searchResultRepository.findByQuery_IdOrderByRankNo(queryId).stream()
.map(VectorSearchCandidate::from)
.toList();
}

private String buildExtractiveFallbackAnswer(List<VectorSearchCandidate> candidates) {
Expand Down
Loading