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
@@ -0,0 +1,5 @@
ALTER TABLE analysis_async_tasks
ADD COLUMN IF NOT EXISTS execution_context_snapshot TEXT;

ALTER TABLE analysis_async_tasks
ADD COLUMN IF NOT EXISTS input_fingerprint_snapshot VARCHAR(64);
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,43 @@ public record AnalysisWorkerContextResponse(
String bigClassificationName,
String middleClassificationName,
String detailClassificationName,
List<AnalysisWorkerQuestionItem> questions
List<AnalysisWorkerQuestionItem> questions,
List<SimilarJobPostingContext> similarJobPostings
) {
public AnalysisWorkerContextResponse(
Long userId,
Long mockApplyId,
String companyName,
String jobTitle,
String task,
String requirements,
String preferredQualifications,
String bigClassificationName,
String middleClassificationName,
String detailClassificationName,
List<AnalysisWorkerQuestionItem> questions
) {
this(
userId,
mockApplyId,
companyName,
jobTitle,
task,
requirements,
preferredQualifications,
bigClassificationName,
middleClassificationName,
detailClassificationName,
questions,
List.of()
);
}

public AnalysisWorkerContextResponse {
questions = questions == null ? List.of() : List.copyOf(questions);
similarJobPostings = similarJobPostings == null ? List.of() : List.copyOf(similarJobPostings);
}

public record AnalysisWorkerQuestionItem(
Long questionId,
String question,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
package com.jobdri.jobdri_api.domain.analysis.dto.worker;

public record SimilarJobPostingContext(
Long jobPostingId,
String companyName,
String postingName,
String jobTitle,
String task,
String requirements,
String preferredQualifications,
int similarityRank,
double similarityScore
) {
}
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,12 @@ public class AnalysisAsyncTask extends CreatedAtEntity {
@Column(name = "estimated_remaining_seconds")
private Integer estimatedRemainingSeconds;

@Column(name = "execution_context_snapshot", columnDefinition = "TEXT")
private String executionContextSnapshot;

@Column(name = "input_fingerprint_snapshot", length = 64)
private String inputFingerprintSnapshot;

public static AnalysisAsyncTask pending(Long userId, Long mockApplyId, int maxRetryCount) {
AnalysisAsyncTask task = new AnalysisAsyncTask();
task.taskId = UUID.randomUUID().toString();
Expand Down Expand Up @@ -217,6 +223,14 @@ public void updateWorkerMetadata(String workerId, Long queueLatencyMillis) {
}
}

public void captureExecutionSnapshot(String executionContextSnapshot, String inputFingerprintSnapshot) {
if (this.executionContextSnapshot != null || this.inputFingerprintSnapshot != null) {
return;
}
this.executionContextSnapshot = executionContextSnapshot;
this.inputFingerprintSnapshot = inputFingerprintSnapshot;
}

private boolean isTerminal() {
return status == TaskStatus.SUCCEEDED || status == TaskStatus.FAILED || status == TaskStatus.CANCELLED;
}
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
package com.jobdri.jobdri_api.domain.analysis.service.async;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.jobdri.jobdri_api.domain.analysis.dto.llm.AnalysisLlmResponse;
import com.jobdri.jobdri_api.domain.analysis.dto.response.AnalysisResponse;
import com.jobdri.jobdri_api.domain.analysis.dto.worker.AnalysisWorkerCompleteRequest;
Expand All @@ -12,6 +14,7 @@
import com.jobdri.jobdri_api.domain.analysis.entity.Question;
import com.jobdri.jobdri_api.domain.analysis.repository.AnalysisAsyncTaskRepository;
import com.jobdri.jobdri_api.domain.analysis.service.core.AnalysisExecutionPayload;
import com.jobdri.jobdri_api.domain.analysis.service.core.AnalysisInputFingerprintProvider;
import com.jobdri.jobdri_api.domain.analysis.service.core.AnalysisService;
import com.jobdri.jobdri_api.domain.user.entity.User;
import com.jobdri.jobdri_api.domain.user.service.UserService;
Expand Down Expand Up @@ -43,6 +46,8 @@ public class AnalysisWorkerBridgeService {
private final AnalysisService analysisService;
private final UserService userService;
private final WorkerTaskResultService workerTaskResultService;
private final AnalysisInputFingerprintProvider analysisInputFingerprintProvider;
private final ObjectMapper objectMapper;

@Transactional
public void markRunning(String taskId, String workerId, int retryCount, Instant submittedAt) {
Expand Down Expand Up @@ -109,10 +114,13 @@ public AnalysisWorkerContextResponse getContext(String taskId, Long userId, Long
);
}
reserveCreditIfNeeded(task);
if (task.getExecutionContextSnapshot() != null) {
return readContextSnapshot(task);
}

User user = userService.getUser(userId);
AnalysisExecutionPayload payload = analysisService.prepareAnalysisExecution(user, mockApplyId);

return new AnalysisWorkerContextResponse(
AnalysisWorkerContextResponse context = new AnalysisWorkerContextResponse(
userId,
mockApplyId,
payload.jobPosting().getCompany().getName(),
Expand All @@ -123,8 +131,14 @@ public AnalysisWorkerContextResponse getContext(String taskId, Long userId, Long
payload.jobPosting().getDetailClassification().getMiddleClassification().getClassification().getBigName(),
payload.jobPosting().getDetailClassification().getMiddleClassification().getMiddleName(),
payload.jobPosting().getDetailClassification().getDetailName(),
toQuestionItems(payload.questions())
toQuestionItems(payload.questions()),
payload.similarJobPostings()
Comment thread
coderabbitai[bot] marked this conversation as resolved.
);
task.captureExecutionSnapshot(
writeContextSnapshot(context),
analysisInputFingerprintProvider.create(payload)
);
return context;
}

@Transactional
Expand Down Expand Up @@ -169,9 +183,20 @@ public AnalysisResponse completeTask(String taskId, AnalysisWorkerCompleteReques
}

User user = userService.getUser(request.userId());
AnalysisExecutionPayload payload = analysisService.prepareAnalysisExecution(user, request.mockApplyId());
AnalysisWorkerContextResponse contextSnapshot = readContextSnapshot(task);
AnalysisExecutionPayload payload = analysisService.prepareAnalysisExecution(
user,
request.mockApplyId(),
contextSnapshot.similarJobPostings()
);
AnalysisLlmResponse llmResponse = request.llmResponse();
AnalysisResponse response = analysisService.finalizeAnalysis(user, request.mockApplyId(), payload, llmResponse);
AnalysisResponse response = analysisService.finalizeAnalysis(
user,
request.mockApplyId(),
payload,
llmResponse,
task.getInputFingerprintSnapshot()
);
analysisAsyncTaskService.updateWorkerMetadata(taskId, request.workerId(), request.queueLatencyMillis());
confirmCreditIfNeeded(task);
analysisAsyncTaskService.markSuccess(taskId, response);
Expand Down Expand Up @@ -219,6 +244,34 @@ private List<AnalysisWorkerContextResponse.AnalysisWorkerQuestionItem> toQuestio
.toList();
}

private String writeContextSnapshot(AnalysisWorkerContextResponse context) {
try {
return objectMapper.writeValueAsString(context);
} catch (JsonProcessingException exception) {
throw new GeneralException(
GeneralErrorCode.INTERNAL_SERVER_ERROR,
"자소서 분석 worker 컨텍스트 snapshot 저장에 실패했습니다."
);
}
}

private AnalysisWorkerContextResponse readContextSnapshot(AnalysisAsyncTask task) {
if (task.getExecutionContextSnapshot() == null || task.getInputFingerprintSnapshot() == null) {
throw new GeneralException(
GeneralErrorCode.INTERNAL_SERVER_ERROR,
"자소서 분석 worker 실행 snapshot이 존재하지 않습니다. taskId=" + task.getTaskId()
);
}
try {
return objectMapper.readValue(task.getExecutionContextSnapshot(), AnalysisWorkerContextResponse.class);
} catch (JsonProcessingException exception) {
throw new GeneralException(
GeneralErrorCode.INTERNAL_SERVER_ERROR,
"자소서 분석 worker 컨텍스트 snapshot을 읽을 수 없습니다. taskId=" + task.getTaskId()
);
}
}

private AnalysisAsyncTask getTask(String taskId) {
return analysisAsyncTaskRepository.findById(taskId)
.orElseThrow(() -> new GeneralException(
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package com.jobdri.jobdri_api.domain.analysis.service.core;

import com.jobdri.jobdri_api.domain.analysis.dto.criteria.JobCategoryEvaluationCriteria;
import com.jobdri.jobdri_api.domain.analysis.dto.worker.SimilarJobPostingContext;
import com.jobdri.jobdri_api.domain.analysis.entity.Question;
import com.jobdri.jobdri_api.domain.corpus.service.CorpusRetrievalService.RetrievalContext;
import com.jobdri.jobdri_api.domain.jobposting.entity.JobPosting;
Expand All @@ -15,7 +16,8 @@ public record AnalysisExecutionPayload(
List<Question> questions,
List<Question> answeredQuestions,
JobCategoryEvaluationCriteria jobCategoryEvaluationCriteria,
RetrievalContext retrievalContext
RetrievalContext retrievalContext,
List<SimilarJobPostingContext> similarJobPostings
) {
public AnalysisExecutionPayload(
Long userId,
Expand All @@ -24,7 +26,7 @@ public AnalysisExecutionPayload(
List<Question> questions,
List<Question> answeredQuestions
) {
this(userId, mockApplyId, jobPosting, questions, answeredQuestions, null, null);
this(userId, mockApplyId, jobPosting, questions, answeredQuestions, null, null, List.of());
}

public AnalysisExecutionPayload(
Expand All @@ -35,6 +37,33 @@ public AnalysisExecutionPayload(
List<Question> answeredQuestions,
JobCategoryEvaluationCriteria jobCategoryEvaluationCriteria
) {
this(userId, mockApplyId, jobPosting, questions, answeredQuestions, jobCategoryEvaluationCriteria, null);
this(userId, mockApplyId, jobPosting, questions, answeredQuestions, jobCategoryEvaluationCriteria, null, List.of());
}

public AnalysisExecutionPayload(
Long userId,
Long mockApplyId,
JobPosting jobPosting,
List<Question> questions,
List<Question> answeredQuestions,
JobCategoryEvaluationCriteria jobCategoryEvaluationCriteria,
RetrievalContext retrievalContext
) {
this(
userId,
mockApplyId,
jobPosting,
questions,
answeredQuestions,
jobCategoryEvaluationCriteria,
retrievalContext,
List.of()
);
}

public AnalysisExecutionPayload {
questions = questions == null ? List.of() : List.copyOf(questions);
answeredQuestions = answeredQuestions == null ? List.of() : List.copyOf(answeredQuestions);
similarJobPostings = similarJobPostings == null ? List.of() : List.copyOf(similarJobPostings);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.jobdri.jobdri_api.domain.analysis.entity.Question;
import com.jobdri.jobdri_api.domain.analysis.dto.worker.SimilarJobPostingContext;
import com.jobdri.jobdri_api.domain.corpus.service.CorpusRetrievalService.RetrievalContext;
import com.jobdri.jobdri_api.domain.corpus.service.CorpusRetrievalService.RetrievedJobPostingReference;
import com.jobdri.jobdri_api.domain.corpus.service.CorpusRetrievalService.RetrievedQuestionReference;
Expand All @@ -25,8 +26,8 @@
@Component
public class AnalysisInputFingerprintProvider {

private static final String FINGERPRINT_SCHEMA_VERSION = "analysis-input-fingerprint-v1";
private static final String ANALYSIS_PROMPT_POLICY_VERSION = "analysis-prompt-policy-v1";
private static final String FINGERPRINT_SCHEMA_VERSION = "analysis-input-fingerprint-v2";
private static final String ANALYSIS_PROMPT_POLICY_VERSION = "analysis-prompt-policy-v2-similar-job-posting-rag";
private static final double ANALYSIS_TEMPERATURE = 0.2;

private final ObjectMapper objectMapper;
Expand Down Expand Up @@ -68,6 +69,7 @@ public String create(AnalysisExecutionPayload payload) {
fingerprintSource.put("fewShotPrompt", fewShotPromptProvider.getPrompt());
fingerprintSource.put("retrievalPolicy", retrievalPolicy());
fingerprintSource.put("retrievalContext", retrievalContextFingerprintSource(payload.retrievalContext()));
fingerprintSource.put("similarJobPostings", similarJobPostingFingerprintSource(payload.similarJobPostings()));
fingerprintSource.put("jobPosting", jobPostingFingerprintSource(payload.jobPosting()));
fingerprintSource.put("answeredQuestions", answeredQuestionFingerprintSource(payload.answeredQuestions()));
fingerprintSource.put("jobCategoryEvaluationCriteria", payload.jobCategoryEvaluationCriteria());
Expand Down Expand Up @@ -149,6 +151,28 @@ private Map<String, Object> jobPostingFingerprintSource(JobPosting jobPosting) {
return jobPostingSource;
}

private List<Map<String, Object>> similarJobPostingFingerprintSource(
List<SimilarJobPostingContext> similarJobPostings
) {
if (similarJobPostings == null) {
return List.of();
}
return similarJobPostings.stream()
.map(context -> {
Map<String, Object> source = new LinkedHashMap<>();
source.put("jobPostingId", context.jobPostingId());
source.put("companyName", defaultString(context.companyName()));
source.put("postingName", defaultString(context.postingName()));
source.put("jobTitle", defaultString(context.jobTitle()));
source.put("task", defaultString(context.task()));
source.put("requirements", defaultString(context.requirements()));
source.put("preferredQualifications", defaultString(context.preferredQualifications()));
source.put("similarityRank", context.similarityRank());
return source;
})
.toList();
}

private List<Map<String, Object>> answeredQuestionFingerprintSource(List<Question> answeredQuestions) {
return answeredQuestions.stream()
.map(question -> {
Expand Down
Loading
Loading