Skip to content
Merged
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,9 @@ public class RagResponse extends BaseEntity {
@JoinColumn(name = "query_id", nullable = false)
private SearchQuery query;

// PROCESSING 상태로 처음 저장될 때는 아직 값이 없다 — Worker가 생성을 마치면 markSuccess/markFailed로 채운다.
// PROCESSING 상태로 처음 저장될 때는 아직 값이 없다 — Worker의 정상 완료(completeSuccess/
// completeFailed) 또는 RagJobTimeoutSweeper의 강제 종료(forceFailIfProcessing) 중 먼저
// 확정되는 쪽이 채운다.
@Column(name = "answer_text", columnDefinition = "TEXT")
private String answerText;

Expand Down Expand Up @@ -100,23 +102,4 @@ 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
Expand Up @@ -52,6 +52,12 @@ public interface RagResponseRepository extends JpaRepository<RagResponse, Long>
* UPDATE로 "이미 끝난 job을 덮어쓰는" 경합을 막는다. 반환값(영향받은 행 수)으로 호출자가
* 실제로 강제 종료가 일어났는지 판단한다.
*
* <p>{@code RagJobTimeoutSweeper}뿐 아니라 {@code RagResponseCommandService.completeFailed()}
* (RagJobWorker가 Ollama 호출 실패를 처리하는 정상 경로)도 이 메서드를 그대로 재사용한다 —
* 둘 다 "PROCESSING인 job을 FAILED + 문구로 확정한다"는 동일한 SQL이 필요하고, 반대로
* RagJobTimeoutSweeper가 먼저 이 job을 확정해버렸다면 RagJobWorker 쪽 시도도 똑같이
* 무시돼야 하기 때문이다(#288).
*
* <p>{@code clearAutomatically}: 벌크 UPDATE는 영속성 컨텍스트를 거치지 않고 DB에 직접
* 실행되므로, 같은 트랜잭션에서 이 job 엔티티를 이미 로딩해둔 상태라면 그 캐시된 인스턴스가
* 여전히 갱신 전 값을 들고 있다 — 이후 같은 트랜잭션에서 다시 조회해도 DB가 아니라 그 캐시를
Expand All @@ -64,4 +70,23 @@ public interface RagResponseRepository extends JpaRepository<RagResponse, Long>
+ "WHERE r.id = :id AND r.status = com.opensource.docgrid.domain.search.enums.ResultStatus.PROCESSING")
int forceFailIfProcessing(@Param("id") Long id, @Param("answerText") String answerText,
@Param("errorMessage") String errorMessage);

/**
* PROCESSING 상태인 job을 SUCCESS + 생성 결과로 확정한다. {@link #forceFailIfProcessing}과
* 대칭되는 목적이다 — RagJobTimeoutSweeper가 이 job을 먼저 FAILED로 강제 종료했다면,
* RagJobWorker의 뒤늦은 정상 완료 시도가 그 결과를 조건 없이 덮어써버리는 경합(#288)을
* 막는다. {@code WHERE ... AND status = PROCESSING} 조건 덕분에, 스위퍼가 먼저 확정해
* 이 UPDATE 시점에 status가 이미 FAILED라면 영향받은 행이 0건이 된다.
*/
@Modifying(clearAutomatically = true)
@Query("UPDATE RagResponse r SET r.status = com.opensource.docgrid.domain.search.enums.ResultStatus.SUCCESS, "
+ "r.answerText = :answerText, r.llmModelName = :llmModelName, "
+ "r.inputTokenCount = :inputTokenCount, r.outputTokenCount = :outputTokenCount, "
+ "r.latencyMs = :latencyMs "
+ "WHERE r.id = :id AND r.status = com.opensource.docgrid.domain.search.enums.ResultStatus.PROCESSING")
int completeSuccessIfProcessing(@Param("id") Long id, @Param("answerText") String answerText,
@Param("llmModelName") String llmModelName,
@Param("inputTokenCount") Integer inputTokenCount,
@Param("outputTokenCount") Integer outputTokenCount,
@Param("latencyMs") Integer latencyMs);
}
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
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.enums.ResultStatus;
import com.opensource.docgrid.domain.search.repository.SearchResultRepository;
import com.opensource.docgrid.global.exception.DocGridException;
import com.opensource.docgrid.global.exception.ErrorCode;
Expand Down Expand Up @@ -117,15 +118,32 @@ public RagEnqueueOutcome enqueue(Long queryId, String queryText, List<VectorSear
*
* <p>{@code job} 객체가 아니라 {@code jobId}만 받아 이 메서드 자신의 트랜잭션 안에서 다시
* 조회하는 이유: RagJobWorker가 {@code findFirstByStatusOrderByCreatedAtAsc()}로 꺼낸
* job은 그 조회 시점에 트랜잭션이 끝나 detached 상태다. 그 인스턴스를 그대로 받아
* markSuccess/markFailed로 값을 바꿔도 이 메서드의 새 트랜잭션에서는 dirty checking이
* 감지하지 못해 DB에 반영되지 않는다(영원히 PROCESSING으로 남아 Worker가 같은 job을
* 계속 재처리하는 버그로 실제 이어졌었다 — #218). {@code findById(jobId)}로 다시 조회해야
* 반드시 managed 상태로 확보된다.
* job은 그 조회 시점에 트랜잭션이 끝나 detached 상태다. 원래(#218) 이 detached 인스턴스를
* 그대로 받아 필드만 바꾸면 dirty checking이 감지 못해 DB에 반영되지 않는 버그가 있었는데,
* 지금은 완료 처리 자체가 dirty checking에 의존하지 않는다({@link
* RagResponseCommandService#completeSuccess}/{@link RagResponseCommandService#completeFailed}
* 참고, #288) — 그래도 {@code findById(jobId)}로 다시 조회해 이 트랜잭션 시점 기준
* 최신 상태(예: {@code promptText})를 읽는다.
*
* <p>{@code completeSuccess}/{@code completeFailed}는 RagJobTimeoutSweeper가 이미 이
* job을 먼저 확정해버렸으면 {@code false}를 반환한다(#288) — 이 경우 citation 저장을
* 건너뛰고 {@code false}를 그대로 반환해, 호출자(RagJobWorker)가 중복 알림을 보내지
* 않게 한다.
*
* <p>{@code findById} 직후 status가 이미 PROCESSING이 아니면 곧바로 {@code false}를
* 반환하고 Ollama를 아예 호출하지 않는다 — RagJobWorker가 이 job을 집어든 뒤, 여기서
* {@code findById}로 다시 읽기 전에 RagJobTimeoutSweeper가 먼저 강제 종료했을 수 있다.
* 이 조기 반환이 없으면 이미 끝난 job에도 Ollama 호출(수십 초)을 그대로 낭비하게 되는데,
* Worker/GPU가 1개뿐이라 그 시간만큼 뒤에 대기 중인 다른 job까지 더 늦어진다 — 아래
* completeSuccess/completeFailed의 조건부 UPDATE는 이 조기 체크 "이후"에 벌어지는 경합(더
* 좁은 창)까지 막아주는 최종 방어선이다.
*/
public void processJob(Long jobId) {
public boolean processJob(Long jobId) {
Comment thread
coderabbitai[bot] marked this conversation as resolved.
RagResponse job = ragResponseRepository.findById(jobId)
.orElseThrow(() -> new DocGridException(ErrorCode.RAG_ANSWER_NOT_FOUND));
if (job.getStatus() != ResultStatus.PROCESSING) {
return false;
}
Long queryId = job.getQuery().getId();

OllamaGenerateResult result;
Expand All @@ -138,9 +156,9 @@ public void processJob(Long jobId) {
List<VectorSearchCandidate> candidates = loadCandidates(queryId);
String fallbackAnswer = candidates.isEmpty() ? e.getErrorCode().getMessage()
: buildExtractiveFallbackAnswer(candidates);
ragResponseCommandService.completeFailed(job, fallbackAnswer, e.getMessage());
boolean completed = ragResponseCommandService.completeFailed(job, fallbackAnswer, e.getMessage());
log.warn("[RAG] fallback queryId={} errorCode={}", queryId, e.getErrorCode().getCode());
return;
return completed;
}

// LLM이 무관하다고 판단해 안내 문구로만 답했으면, 근거 문서를 같이 보여주지 않는다. 단, 7B
Expand All @@ -163,15 +181,23 @@ public void processJob(Long jobId) {
}
}

ragResponseCommandService.completeSuccess(job, new OllamaGenerateResult(
boolean completed = ragResponseCommandService.completeSuccess(job, new OllamaGenerateResult(
result.model(), answerText, result.inputTokenCount(), result.outputTokenCount(), result.latencyMs()
));
if (!completed) {
// RagJobTimeoutSweeper가 이 job을 이미 FAILED로 강제 종료한 뒤라는 뜻이다 — 방금
// 만든 답변은 이미 아무도 안 볼 결과라, citation 저장도 하지 않고 그대로 물러난다.
log.info("[RAG] job이 이미 timeout으로 종료됨(경합), 완료 결과 반영 안 함 queryId={} responseId={}",
queryId, job.getId());
return false;
}

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());
return true;
}

/**
Expand All @@ -180,10 +206,15 @@ public void processJob(Long jobId) {
* PROCESSING으로 남아, 같은 job을 Worker가 계속 다시 집어 무한 재시도하게 된다 —
* detached entity 버그(#218)와 증상이 같아진다. {@code ifPresent}로 감싸는 이유는 job이
* 이미 다른 이유로 없어졌을 수 있는 극단적 상황을 방어하기 위함이다.
*
* @return 실제로 이 호출로 FAILED 확정이 일어났으면 true. job이 없거나(극단적 상황),
* RagJobTimeoutSweeper가 이미 먼저 확정해뒀으면(#288) false — 호출자
* (RagJobWorker)는 이 경우 WebSocket 알림을 보내지 않는다.
*/
public void markUnexpectedFailure(Long jobId, String errorMessage) {
ragResponseRepository.findById(jobId)
.ifPresent(job -> ragResponseCommandService.completeFailed(job, UNEXPECTED_FAILURE_ANSWER_TEXT, errorMessage));
public boolean markUnexpectedFailure(Long jobId, String errorMessage) {
return ragResponseRepository.findById(jobId)
.map(job -> ragResponseCommandService.completeFailed(job, UNEXPECTED_FAILURE_ANSWER_TEXT, errorMessage))
.orElse(false);
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,11 @@ public class RagJobWorker {
* 이미 올바르게 반영된 결과를 덮어쓰지 않도록 조용히 넘어간다(알림도 안 보낸다 — 그건
* 먼저 처리한 쪽의 몫). ③그 외 예상 못한 예외 — FAILED로 강제 확정한 뒤 알림까지 보낸다
* (실패했어도 화면이 영원히 로딩중으로 안 남도록).
*
* <p>①/③ 모두 알림은 {@code processJob}/{@code markUnexpectedFailure}가 반환하는
* boolean을 확인한 뒤에만 보낸다 — RagJobTimeoutSweeper가 이 job을 이미 먼저 FAILED로
* 확정해뒀다면(#288) 두 메서드 다 실제로는 아무것도 안 바꾸고 false를 반환하는데, 이 경우
* 스위퍼가 이미 보낸 알림 외에 Worker가 중복으로 또 보낼 이유가 없다.
*/
@Scheduled(fixedDelayString = "${rag.worker.polling-interval:1s}")
public void processNext() {
Expand All @@ -57,14 +62,17 @@ public void processNext() {
RagResponse job = maybeJob.get();
/*
* query/query.user는 findFirstByStatusOrderByCreatedAtAsc()의 @EntityGraph로 이미
* 로딩돼 있어 detached 상태에서 읽어도 안전하다 — 문제는 "쓰기"(markSuccess 등)뿐이라
* processJob()에는 id만 넘겨 그 안에서 managed 상태로 다시 조회하게 한다.
* 로딩돼 있어 detached 상태에서 읽어도 안전하다 — 완료 처리 자체는 조건부 UPDATE로
* 이뤄지므로(#288) 이 job 인스턴스가 detached여도 상관없지만, processJob()이 이
* 트랜잭션 시점 기준 최신 상태를 읽도록 id만 넘긴다.
*/
Long queryId = job.getQuery().getId();
String userEmail = job.getQuery().getUser().getEmail();

try {
ragFacade.processJob(job.getId());
if (ragFacade.processJob(job.getId())) {
ragWebSocketController.notifyAnswerReady(userEmail, queryId);
}
} catch (OptimisticLockingFailureException e) {
/*
* 설계상 Worker는 인스턴스 1개를 전제하지만(클래스 주석 참고), 롤링 배포로 신·구
Expand All @@ -74,7 +82,6 @@ public void processNext() {
* 사고가 난다 — 조용히 다음 폴링으로 넘어간다.
*/
log.warn("[RAG-WORKER] job이 이미 다른 트랜잭션에서 처리된 것으로 보임(경합) queryId={}", queryId);
return;
} catch (Exception e) {
/*
* processJob() 내부에서 Ollama 관련 실패는 이미 DocGridException으로 잡아 fallback
Expand All @@ -84,11 +91,9 @@ public void processNext() {
* 되므로(detached entity 버그와 같은 증상), 반드시 FAILED로 확정한 뒤 넘어간다.
*/
log.error("[RAG-WORKER] job 처리 중 예상치 못한 예외 queryId={}", queryId, e);
ragFacade.markUnexpectedFailure(job.getId(), e.getMessage());
ragWebSocketController.notifyAnswerReady(userEmail, queryId);
return;
if (ragFacade.markUnexpectedFailure(job.getId(), e.getMessage())) {
ragWebSocketController.notifyAnswerReady(userEmail, queryId);
}
}

ragWebSocketController.notifyAnswerReady(userEmail, queryId);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,10 @@
* RAG 최종 답변 저장 서비스 (F-RAG-03).
*
* <p>비동기 Job 큐 전환(#218) 이후에는 검색 직후 PROCESSING row를 먼저 저장해두고(createPending),
* Worker가 LLM 생성을 마친 뒤 completeSuccess/completeFailed로 같은 row를 채운다(내부적으로
* RagResponse 엔티티의 markSuccess/markFailed를 호출) — SearchQuery의 PROCESSING 선저장 패턴과
* 동일해졌다.
* Worker가 LLM 생성을 마친 뒤 completeSuccess/completeFailed로 같은 row를 채운다 — SearchQuery의
* PROCESSING 선저장 패턴과 동일해졌다. 이 둘은 RagJobTimeoutSweeper와의 경합(#288) 때문에
* 엔티티 dirty checking이 아니라 {@code RagResponseRepository}의 조건부 UPDATE(WHERE
* status=PROCESSING)로 직접 확정하고, 실제로 확정이 일어났는지를 boolean으로 반환한다.
*/
@Transactional
@Service
Expand Down Expand Up @@ -60,27 +61,46 @@ public RagResponse createNoContext(SearchQuery query) {

/**
* Worker가 LLM 생성에 성공했을 때, createPending()으로 미리 저장해둔 row를 SUCCESS로
* 채운다. ragResponse는 이미 영속 상태라 save()를 다시 부르지 않아도 트랜잭션 커밋 시점에
* 더티체킹으로 자동 반영된다.
* 채운다.
*
* <p>엔티티를 불러와 마크하고 dirty checking에 맡기는 대신, {@link
* RagResponseRepository#completeSuccessIfProcessing}(조건부 UPDATE, {@code WHERE
* status = PROCESSING})으로 직접 확정한다 — RagJobTimeoutSweeper가 이 job을 먼저 FAILED로
* 강제 종료했다면, Worker의 이 뒤늦은 성공 처리가 그 결과를 조건 없이 덮어써버리는 경합
* (#288)을 막기 위함이다. 영향받은 행이 0건이면(=스위퍼가 먼저 확정함) {@code false}를
* 반환하고, 호출자(RagFacade.processJob)는 이 경우 citation 저장도 건너뛴다 — 이미 아무도
* 안 볼 결과이기 때문이다.
*
* @return 실제로 이 호출로 SUCCESS 확정이 일어났으면 true, 이미 다른 경로(스위퍼)가
* 먼저 끝내 아무 일도 하지 않았으면 false.
*/
public void completeSuccess(RagResponse ragResponse, OllamaGenerateResult result) {
ragResponse.markSuccess(
result.answerText(), result.model(), result.inputTokenCount(), result.outputTokenCount(),
result.latencyMs()
public boolean completeSuccess(RagResponse ragResponse, OllamaGenerateResult result) {
int updated = ragResponseRepository.completeSuccessIfProcessing(
ragResponse.getId(), result.answerText(), result.model(),
result.inputTokenCount(), result.outputTokenCount(), result.latencyMs()
);
return updated > 0;
}

/**
* Worker가 LLM 생성에 실패했을 때, 빈손 대신 fallbackAnswerText(대개 extractive fallback)를
* 채우고 status만 FAILED로 남긴다.
*
* <p>{@link RagResponseRepository#forceFailIfProcessing}(조건부 UPDATE)을
* RagJobTimeoutSweeper와 공유해서 쓴다 — 이유는 {@link #completeSuccess}와 동일하다.
*
* <p>검색 도메인의 SearchQueryCommandService.markFailed()와 달리 REQUIRES_NEW가 없다 —
* 이 메서드를 부르는 RagFacade.processJob()의 catch 블록은 예외를 다시 던지지 않고 그대로
* return하므로, 이 메서드가 실행되는 트랜잭션 자체가 롤백될 일이 없다. 재전파해서 바깥
* 트랜잭션을 일부러 굴리는 검색 쪽 구조와 달리, 애초에 롤백될 트랜잭션이 없어 REQUIRES_NEW로
* 실패 기록을 따로 지킬 필요 자체가 없다.
*
* @return 실제로 이 호출로 FAILED 확정이 일어났으면 true, 이미 다른 경로(스위퍼)가
* 먼저 끝내 아무 일도 하지 않았으면 false.
*/
public void completeFailed(RagResponse ragResponse, String fallbackAnswerText, String errorMessage) {
ragResponse.markFailed(fallbackAnswerText, errorMessage);
public boolean completeFailed(RagResponse ragResponse, String fallbackAnswerText, String errorMessage) {
int updated = ragResponseRepository.forceFailIfProcessing(
ragResponse.getId(), fallbackAnswerText, errorMessage);
return updated > 0;
}
}
Loading