Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
package repit.repit_api_server.domain.metadata.controller;

import lombok.RequiredArgsConstructor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;
Expand All @@ -18,6 +20,7 @@
import repit.repit_api_server.global.client.AiServerClient;
import repit.repit_api_server.global.client.AuthServerClient;
import repit.repit_api_server.global.common.ApiResponse;
import repit.repit_api_server.global.exception.ExternalApiException;
import repit.repit_api_server.global.response.UserResponse;

import java.io.IOException;
Expand All @@ -27,6 +30,9 @@
@RequiredArgsConstructor
@RequestMapping("/api/v1/ai")
public class AiMetaDataController {

private static final Logger log = LoggerFactory.getLogger(AiMetaDataController.class);

private static final long SSE_TIMEOUT = 10 * 60 * 1000L; // 10분

@Value("${app.callback-base-url}")
Expand All @@ -42,6 +48,16 @@ public class AiMetaDataController {
@GetMapping("/subscribe/{jobId}")
public SseEmitter subscribe(@PathVariable String jobId) {
SseEmitter emitter = new SseEmitter(SSE_TIMEOUT);

// 콜백이 구독보다 먼저 도착했을 수 있다. 그때는 붙는 즉시 결과를 돌려주고 끝낸다.
// 되짚어주지 않으면 이미 끝난 작업을 구독한 클라이언트는 아무것도 받지 못한 채 타임아웃까지
// 매달려 있고, EventSource가 그때마다 다시 붙어 재연결만 반복한다.
CallbackSuccessResponse finished = aiMetaDataService.findFinished(jobId);
if (finished != null) {
sendCompletionEvent(emitter, jobId, finished);
return emitter;
}

sseEmitterRepository.save(jobId, emitter);

emitter.onCompletion(() -> sseEmitterRepository.remove(jobId));
Expand Down Expand Up @@ -103,13 +119,35 @@ public ResponseEntity<GenerateResponse> generateMock(
return ResponseEntity.ok(response);
}

// 이후 채팅 서버 요청에서 jobId를 서버가 직접 찾을 수 있도록 소유자를 기록해둔다.
/**
* 이후 채팅 서버 요청에서 jobId를 서버가 직접 찾을 수 있도록 소유자를 기록해둔다.
*
* <p>소유자를 남기지 못하면 이 분석 결과는 사용자로 되찾을 수 없어 면접 질문 재작성이
* 예전 결과를 집어 든다. 그래서 실패를 조용히 넘기지 않고 반드시 로그로 남긴다.
*
* <p>다만 이 시점에는 분석 서버가 이미 작업을 접수한 뒤다. 소유자 기록이 실패했다고 요청
* 전체를 실패시키면 클라이언트가 jobId를 받지 못해 결과를 영영 조회할 수 없게 되므로,
* 기록 실패는 예외로 번지지 않게 막는다.
*/
private void registerJobOwner(String authorization, GenerateResponse response) {
UserResponse user = authServerClient.getUser(authorization);
if (user == null) {
return;
if (response == null || response.getJob_id() == null) {
// jobId가 없으면 구독도 조회도 할 수 없다. 성공으로 돌려주면 원인을 찾을 수 없다.
log.error("분석 서버 응답에 job_id가 없습니다. status={}, message={}",
response == null ? null : response.getStatus(),
response == null ? null : response.getMessage());
throw new ExternalApiException("분석 서버가 작업 번호를 돌려주지 않았습니다.", null, null);
}

try {
UserResponse user = authServerClient.getUser(authorization);
if (user == null || user.getId() == null) {
log.error("분석 작업의 소유자를 확인하지 못했습니다. jobId={}", response.getJob_id());
return;
}
aiMetaDataService.registerJob(response.getJob_id(), user.getId());
} catch (RuntimeException e) {
log.error("분석 작업의 소유자를 기록하지 못했습니다. jobId={}", response.getJob_id(), e);
}
aiMetaDataService.registerJob(response.getJob_id(), user.getId());
}

@PostMapping("/callback")
Expand All @@ -132,9 +170,13 @@ public ApiResponse<CallbackSuccessResponse> callback(
private void sendCompletionEvent(String jobId, CallbackSuccessResponse response) {
SseEmitter emitter = sseEmitterRepository.get(jobId);
if (emitter == null) {
// 아직 아무도 구독하지 않았다. 결과는 DB에 있으니 구독이 붙을 때 되짚어 보낸다.
return;
}
sendCompletionEvent(emitter, jobId, response);
}

private void sendCompletionEvent(SseEmitter emitter, String jobId, CallbackSuccessResponse response) {
String eventName = "succeeded".equalsIgnoreCase(response.getStatus())
? "question-generated"
: "question-generation-failed";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@

import jakarta.persistence.Column;
import jakarta.persistence.Entity;
import jakarta.persistence.EnumType;
import jakarta.persistence.Enumerated;
import jakarta.persistence.Id;
import jakarta.persistence.Table;
import lombok.AllArgsConstructor;
Expand All @@ -12,6 +14,7 @@
import org.hibernate.annotations.CreationTimestamp;
import org.hibernate.annotations.JdbcTypeCode;
import org.hibernate.type.SqlTypes;
import repit.repit_api_server.domain.metadata.entity.enums.AnalysisStatus;

import java.time.LocalDateTime;

Expand All @@ -30,10 +33,22 @@ public class AnalysisDataEntity {
@Column(name = "user_id")
private Long userId;

// 콜백이 오기 전에는 PENDING이다. 실패한 작업과 아직 끝나지 않은 작업을 구분하려면 이 값이 필요하다.
@Enumerated(EnumType.STRING)
@Column(nullable = false)
@Builder.Default
private AnalysisStatus status = AnalysisStatus.PENDING;

@JdbcTypeCode(SqlTypes.JSON)
@Column(columnDefinition = "jsonb")
private Object result;

// 실패 콜백에만 채워진다.
private Integer errorStatusCode;

@Column(columnDefinition = "TEXT")
private String errorMessage;

@CreationTimestamp
@Column(nullable = false, updatable = false)
private LocalDateTime createdAt;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
package repit.repit_api_server.domain.metadata.entity.enums;

/** 분석 작업의 상태. 콜백이 오기 전까지는 PENDING이다. */
public enum AnalysisStatus {
PENDING,
SUCCEEDED,
FAILED
}
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
package repit.repit_api_server.domain.metadata.repository;

import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.data.jpa.repository.Modifying;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param;
import repit.repit_api_server.domain.metadata.entity.AnalysisDataEntity;

import java.util.Optional;
Expand All @@ -9,4 +12,9 @@ public interface AnalysisDataRepository extends JpaRepository<AnalysisDataEntity

// 분석이 끝난(result가 채워진) 가장 최근 작업
Optional<AnalysisDataEntity> findTopByUserIdAndResultIsNotNullOrderByCreatedAtDesc(Long userId);

// 소유자만 갱신한다. 엔티티를 통째로 저장하면 콜백이 먼저 채워둔 result를 덮어쓸 수 있다.
@Modifying(clearAutomatically = true, flushAutomatically = true)
@Query("update AnalysisDataEntity a set a.userId = :userId where a.jobId = :jobId")
int updateUserId(@Param("jobId") String jobId, @Param("userId") Long userId);
}
Original file line number Diff line number Diff line change
@@ -1,39 +1,128 @@
package repit.repit_api_server.domain.metadata.service;

import lombok.RequiredArgsConstructor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import repit.repit_api_server.domain.metadata.dto.request.CallbackSuccessRequest;
import repit.repit_api_server.domain.metadata.dto.response.CallbackSuccessResponse;
import repit.repit_api_server.domain.metadata.dto.response.ResultResponse;
import repit.repit_api_server.domain.metadata.entity.AnalysisDataEntity;
import repit.repit_api_server.domain.metadata.entity.enums.AnalysisStatus;
import repit.repit_api_server.domain.metadata.repository.AnalysisDataRepository;

@Service
@RequiredArgsConstructor
public class AiMetaDataService {

private static final Logger log = LoggerFactory.getLogger(AiMetaDataService.class);

private static final String STATUS_SUCCEEDED = "succeeded";
private static final String STATUS_FAILED = "failed";

private final AnalysisDataRepository analysisDataRepository;

// 분석 요청 시점에 작업 소유자를 먼저 기록해둔다. 결과는 콜백에서 채워진다.
/**
* 분석 요청 시점에 작업 소유자를 먼저 기록해둔다. 결과는 콜백에서 채워진다.
*
* <p>분석 서버가 빠르면 이 메서드보다 콜백이 먼저 도착할 수 있다. 그때 행 전체를 저장하면
* 아직 비어 있는 result로 이미 받아둔 결과를 덮어쓰게 되므로, 소유자 컬럼만 갱신하고
* 행이 없을 때만 새로 만든다.
*/
@Transactional
public void registerJob(String jobId, Long userId) {
if (jobId == null || userId == null) {
return;
}
AnalysisDataEntity data = analysisDataRepository.findById(jobId)
.orElseGet(() -> AnalysisDataEntity.builder().jobId(jobId).build());
data.setUserId(userId);
analysisDataRepository.save(data);
if (analysisDataRepository.updateUserId(jobId, userId) == 0) {
analysisDataRepository.save(AnalysisDataEntity.builder()
.jobId(jobId)
.userId(userId)
.build());
}
}

/**
* 분석 결과 콜백을 저장한다. 재전송이 있을 수 있어 두 번 받아도 안전해야 한다.
*
* <p>실패 콜백에는 result가 없다. status를 보지 않고 그대로 덮어쓰면 먼저 받아둔 성공
* 결과까지 지워지므로, 성공 콜백일 때만 결과를 저장한다.
*/
@Transactional
public void saveResult(CallbackSuccessRequest request) {
String jobId = request.getJob_id();
if (jobId == null) {
log.warn("job_id 없는 분석 콜백을 받았습니다. status={}", request.getStatus());
return;
}

// registerJob으로 이미 저장된 행이 있으면 userId를 유지한 채 결과만 채운다.
AnalysisDataEntity data = analysisDataRepository.findById(request.getJob_id())
.orElseGet(() -> AnalysisDataEntity.builder().jobId(request.getJob_id()).build());
data.setResult(request.getResult());
AnalysisDataEntity data = analysisDataRepository.findById(jobId)
.orElseGet(() -> AnalysisDataEntity.builder().jobId(jobId).build());

if (STATUS_SUCCEEDED.equalsIgnoreCase(request.getStatus()) && request.getResult() != null) {
data.setStatus(AnalysisStatus.SUCCEEDED);
data.setResult(request.getResult());
data.setErrorStatusCode(null);
data.setErrorMessage(null);
analysisDataRepository.save(data);
return;
}

// 이미 결과를 받아둔 작업이라면 뒤늦은 실패 콜백에 그 결과를 잃을 이유가 없다.
if (data.getStatus() == AnalysisStatus.SUCCEEDED) {
log.warn("이미 성공한 분석에 실패 콜백이 도착해 무시합니다. jobId={}, error={}",
jobId, describeError(request.getError()));
return;
}

log.warn("분석에 실패했습니다. jobId={}, status={}, error={}",
jobId, request.getStatus(), describeError(request.getError()));
data.setStatus(AnalysisStatus.FAILED);
if (request.getError() != null) {
data.setErrorStatusCode(request.getError().getStatus_code());
data.setErrorMessage(request.getError().getMessage());
}
analysisDataRepository.save(data);
}

/**
* 이미 끝난 작업이면 콜백과 같은 모양으로 돌려준다. 아직 진행 중이거나 모르는 작업이면 null.
*
* <p>구독이 콜백보다 늦게 붙는 경우가 있다. 그때 SSE로 흘릴 것이 없으면 클라이언트는
* 영영 아무것도 받지 못하므로, 저장해둔 결과로 되짚어준다.
*/
@Transactional(readOnly = true)
public CallbackSuccessResponse findFinished(String jobId) {
return analysisDataRepository.findById(jobId)
.filter(data -> data.getStatus() != AnalysisStatus.PENDING)
.map(data -> CallbackSuccessResponse.builder()
.job_id(data.getJobId())
.status(data.getStatus() == AnalysisStatus.SUCCEEDED ? STATUS_SUCCEEDED : STATUS_FAILED)
.result(data.getResult())
.error(toError(data))
.build())
.orElse(null);
}

private CallbackSuccessRequest.Error toError(AnalysisDataEntity data) {
if (data.getErrorStatusCode() == null && data.getErrorMessage() == null) {
return null;
}
return CallbackSuccessRequest.Error.builder()
.status_code(data.getErrorStatusCode())
.message(data.getErrorMessage())
.build();
}

private String describeError(CallbackSuccessRequest.Error error) {
if (error == null) {
return null;
}
return error.getStatus_code() + " " + error.getMessage();
}

/**
* 저장된 분석 결과를 원형 그대로 돌려준다.
* 면접에 쓸 질문은 재작성이 끝나는 시점에 이 서버가 채팅 서버로 직접 넘기므로,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,8 @@ public class AnswerEntity {
@Column(nullable = false)
private int responseTime;

@Column(nullable = false)
// 모의면접 답변은 길다. 255자로 자르면 그대로 피드백 품질이 깎인다.
@Column(nullable = false, columnDefinition = "TEXT")
private String content;

@CreationTimestamp
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,8 @@ public class FeedbackCallbackRequest {
@AllArgsConstructor
public static class Result {
private Overall overall;
// N:1 면접에만 실린다. 1:1 콜백에는 없다.
private List<Persona> personas;
private List<Item> feedbacks;
}

Expand All @@ -43,11 +45,28 @@ public static class Overall {
private Integer questionCount;
}

/** 면접관별 종합. 문항이 2~3개뿐이라 점수는 하나만 온다. */
@Getter
@NoArgsConstructor
@AllArgsConstructor
public static class Persona {
private Long personaId;
private String personaRole;
private Integer score;
private String comment;
private List<String> strengths;
private List<String> improvements;
private Integer answeredCount;
private Integer questionCount;
}

@Getter
@NoArgsConstructor
@AllArgsConstructor
public static class Item {
private String questionId;
// N:1에서 이 질문을 던진 면접관. 1:1 콜백에는 없다.
private Long personaId;
private String questionContent;
private String intention;
private String userAnswer;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
@AllArgsConstructor
public class FeedbackItemResponse {
private String questionId;
private Long personaId;
private String questionContent;
private String intention;
private String userAnswer;
Expand All @@ -25,6 +26,7 @@ public class FeedbackItemResponse {
public static FeedbackItemResponse from(FeedbackItemEntity item) {
return FeedbackItemResponse.builder()
.questionId(item.getQuestionId())
.personaId(item.getPersonaId())
.questionContent(item.getQuestionContent())
.intention(item.getIntention())
.userAnswer(item.getUserAnswer())
Expand Down
Loading