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
Expand Up @@ -15,6 +15,8 @@
import org.hibernate.exception.ConstraintViolationException;
import org.springframework.dao.DataIntegrityViolationException;

import java.util.Optional;

// 분석 비동기 작업의 접수/조회/취소 유스케이스를 조율한다.
public class AnalysisAsyncUseCase {
private static final String ACTIVE_TASK_UNIQUE_CONSTRAINT = "uk_analysis_async_tasks_active_user_mock_apply";
Expand Down Expand Up @@ -64,6 +66,11 @@ public AnalysisAsyncCancelResponse cancel(User user, Long mockApplyId, String ta
}

private AnalysisAsyncSubmitResponse createCachedOrProcessTask(User user, Long mockApplyId) {
Optional<AnalysisAsyncTask> recoverableTask =
analysisAsyncTaskService.findRecoverablePublishFailureTask(user.getId(), mockApplyId);
if (recoverableTask.isPresent()) {
return reopenAndProcessTask(user, mockApplyId, recoverableTask.get());
}
if (analysisService.hasReusableAnalysis(user, mockApplyId)) {
return toCachedResponse();
}
Expand All @@ -77,8 +84,20 @@ private AnalysisAsyncSubmitResponse createAndProcessTask(User user, Long mockApp
}

AnalysisAsyncTask task = pendingTaskResult.task();
String taskId = task.getTaskId();
return processTask(user, mockApplyId, task);
}

private AnalysisAsyncSubmitResponse reopenAndProcessTask(User user, Long mockApplyId, AnalysisAsyncTask task) {
AnalysisAsyncTaskService.ReopenPublishFailureResult reopenResult =
analysisAsyncTaskService.reopenPublishFailureTask(task.getTaskId());
if (!reopenResult.reopened()) {
return toInProgressResponse(reopenResult.task());
}
return processTask(user, mockApplyId, reopenResult.task());
}

private AnalysisAsyncSubmitResponse processTask(User user, Long mockApplyId, AnalysisAsyncTask task) {
String taskId = task.getTaskId();
try {
analysisAsyncProcessor.process(
taskId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,9 @@ public class AnalysisAsyncTask extends CreatedAtEntity {
@Column(name = "credit_reference_id", length = 100)
private String creditReferenceId;

@Column(name = "credit_reference_version", nullable = false)
private int creditReferenceVersion;

@Enumerated(EnumType.STRING)
@Column(name = "credit_status", nullable = false, length = 20)
private AnalysisAsyncCreditStatus creditStatus;
Expand Down Expand Up @@ -107,6 +110,7 @@ public static AnalysisAsyncTask pending(Long userId, Long mockApplyId, int maxRe
task.userId = userId;
task.mockApplyId = mockApplyId;
task.creditStatus = AnalysisAsyncCreditStatus.NONE;
task.creditReferenceVersion = 0;
task.status = AnalysisAsyncTaskStatus.PENDING;
task.message = "자소서 분석 비동기 작업이 접수되었습니다.";
task.retryCount = 0;
Expand All @@ -118,17 +122,39 @@ public static AnalysisAsyncTask pending(Long userId, Long mockApplyId, int maxRe
return task;
}

public void markCreditReserved(String creditReferenceId) {
public boolean markCreditReserved(String creditReferenceId) {
if (!canReserveCredit() || creditReferenceId == null || creditReferenceId.isBlank()) {
return false;
}
this.creditReferenceId = creditReferenceId;
this.creditReferenceVersion += 1;
this.creditStatus = AnalysisAsyncCreditStatus.RESERVED;
return true;
}

public void markCreditConfirmed() {
public boolean markCreditConfirmed() {
if (creditStatus != AnalysisAsyncCreditStatus.RESERVED || creditReferenceId == null) {
return false;
}
this.creditStatus = AnalysisAsyncCreditStatus.CONFIRMED;
return true;
}

public void markCreditReleased() {
public boolean markCreditReleased() {
if (creditStatus != AnalysisAsyncCreditStatus.RESERVED || creditReferenceId == null) {
return false;
}
this.creditStatus = AnalysisAsyncCreditStatus.RELEASED;
return true;
}

public boolean canReserveCredit() {
return creditStatus == AnalysisAsyncCreditStatus.NONE
|| creditStatus == AnalysisAsyncCreditStatus.RELEASED;
}

public int nextCreditReferenceVersion() {
return creditReferenceVersion + 1;
}

public void markRunning(String workerId, int retryCount, Instant messageSubmittedAt) {
Expand Down Expand Up @@ -169,6 +195,25 @@ public boolean isRecoverablePublishFailure() {
&& failureReason == AnalysisAsyncFailureReason.PUBLISH_FAILED;
}

public void reopenForRepublish() {
if (!isRecoverablePublishFailure()) {
return;
}
this.status = AnalysisAsyncTaskStatus.PENDING;
this.message = "자소서 분석 비동기 작업이 다시 접수되었습니다.";
this.error = null;
this.failureReason = null;
this.workerId = null;
this.submittedAt = LocalDateTime.now();
this.lastAttemptAt = null;
this.queueLatencyMillis = null;
this.startedAt = null;
this.completedAt = null;
this.currentStep = "VALIDATING_INPUT";
this.progressPercent = 0;
this.estimatedRemainingSeconds = null;
}

public void markRetryScheduled(AnalysisAsyncFailureReason failureReason, String errorMessage, int retryCount) {
if (isTerminal()) {
return;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,61 +1,102 @@
package com.jobdri.jobdri_api.domain.analysis.infrastructure.async;

import com.jobdri.jobdri_api.domain.analysis.entity.AnalysisAsyncTask;
import com.jobdri.jobdri_api.domain.analysis.type.AnalysisAsyncCreditStatus;
import com.jobdri.jobdri_api.domain.analysis.type.AnalysisAsyncFailureReason;
import com.jobdri.jobdri_api.domain.analysis.type.AnalysisAsyncTaskStatus;
import com.jobdri.jobdri_api.domain.analysis.repository.AnalysisAsyncTaskRepository;
import com.jobdri.jobdri_api.domain.analysis.service.async.AnalysisAsyncCreditCoordinator;
import com.jobdri.jobdri_api.domain.analysis.service.async.AnalysisAsyncTaskService;
import com.jobdri.jobdri_api.domain.analysis.service.async.AnalysisQueueProperties;
import com.jobdri.jobdri_api.domain.analysis.service.core.AnalysisCreditService;
import com.jobdri.jobdri_api.domain.user.entity.User;
import com.jobdri.jobdri_api.domain.user.service.UserService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.domain.PageRequest;
import org.springframework.transaction.support.TransactionTemplate;

import java.time.Clock;
import java.time.Duration;
import java.time.LocalDateTime;
import java.util.EnumSet;
import java.util.List;

@Slf4j
// timeout/retry 기준으로 만료된 분석 async task를 정리한다.
public class AnalysisAsyncTaskSweepCoordinator {
private static final int SWEEP_BATCH_SIZE = 100;

private final AnalysisAsyncTaskRepository analysisAsyncTaskRepository;
private final AnalysisAsyncTaskService analysisAsyncTaskService;
private final AnalysisCreditService analysisCreditService;
private final UserService userService;
private final AnalysisAsyncCreditCoordinator analysisAsyncCreditCoordinator;
private final TransactionTemplate transactionTemplate;
private final AnalysisQueueProperties analysisQueueProperties;
private final Clock clock;

public AnalysisAsyncTaskSweepCoordinator(
AnalysisAsyncTaskRepository analysisAsyncTaskRepository,
AnalysisAsyncTaskService analysisAsyncTaskService,
AnalysisCreditService analysisCreditService,
UserService userService,
AnalysisAsyncCreditCoordinator analysisAsyncCreditCoordinator,
TransactionTemplate transactionTemplate,
AnalysisQueueProperties analysisQueueProperties
AnalysisQueueProperties analysisQueueProperties,
Clock clock
) {
this.analysisAsyncTaskRepository = analysisAsyncTaskRepository;
this.analysisAsyncTaskService = analysisAsyncTaskService;
this.analysisCreditService = analysisCreditService;
this.userService = userService;
this.analysisAsyncCreditCoordinator = analysisAsyncCreditCoordinator;
this.transactionTemplate = transactionTemplate;
this.analysisQueueProperties = analysisQueueProperties;
this.clock = clock;
}

public int sweepTimedOutTasks() {
LocalDateTime now = LocalDateTime.now(clock);
int expiredCount = sweepTimedOutPendingTasks(now);
expiredCount += sweepTimedOutRunningTasks(now);
return expiredCount;
}

private int sweepTimedOutPendingTasks(LocalDateTime now) {
LocalDateTime deadline = now.minusSeconds(analysisQueueProperties.getQueueTimeoutSeconds());
return sweepTimedOutTaskIds(
() -> analysisAsyncTaskRepository.findTimedOutPendingTaskIds(
deadline,
PageRequest.of(0, SWEEP_BATCH_SIZE)
)
);
}

private int sweepTimedOutRunningTasks(LocalDateTime now) {
LocalDateTime deadline = now.minusSeconds(analysisQueueProperties.getProcessingTimeoutSeconds());
return sweepTimedOutTaskIds(
() -> analysisAsyncTaskRepository.findTimedOutRunningTaskIds(
deadline,
PageRequest.of(0, SWEEP_BATCH_SIZE)
)
);
}

private int sweepTimedOutTaskIds(TaskIdBatchLoader taskIdBatchLoader) {
int expiredCount = 0;
for (AnalysisAsyncTask task : analysisAsyncTaskRepository.findByStatusIn(EnumSet.of(AnalysisAsyncTaskStatus.PENDING, AnalysisAsyncTaskStatus.RUNNING))) {
try {
expiredCount += transactionTemplate.execute(status -> sweepTimedOutTask(task.getTaskId()));
} catch (RuntimeException e) {
log.error("Analysis async task sweep failed for taskId={}", task.getTaskId(), e);
while (true) {
List<String> taskIds = taskIdBatchLoader.load();
if (taskIds.isEmpty()) {
return expiredCount;
}
for (String taskId : taskIds) {
expiredCount += sweepTimedOutTask(taskId);
}
if (taskIds.size() < SWEEP_BATCH_SIZE) {
return expiredCount;
}
}
return expiredCount;
}

private int sweepTimedOutTask(String taskId) {
try {
return transactionTemplate.execute(status -> sweepTimedOutTaskInTransaction(taskId));
} catch (RuntimeException e) {
log.error("Analysis async task sweep failed for taskId={}", taskId, e);
return 0;
}
}
Comment on lines +74 to +97

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🩺 Stability & Availability | 🟠 Major | 🏗️ Heavy lift

실패한 첫 batch가 무한 반복될 수 있습니다.

처리 실패 시 task는 PENDING 또는 RUNNING 상태로 남습니다. 다음 반복도 PageRequest.of(0, SWEEP_BATCH_SIZE)를 사용하므로 같은 첫 100개 ID를 다시 조회합니다. 이 상태가 지속되면 sweep 스레드가 종료하지 않습니다.

  • src/main/java/com/jobdri/jobdri_api/domain/analysis/infrastructure/async/AnalysisAsyncTaskSweepCoordinator.java#L80-L103: 실패한 ID를 한 sweep에서 다시 처리하지 않도록 keyset pagination 또는 처리 완료 ID 추적을 적용하세요.
  • src/test/java/com/jobdri/jobdri_api/domain/analysis/infrastructure/async/AnalysisAsyncTaskSweepCoordinatorTest.java#L94-L97: 100개 ID가 모두 실패하는 경우 sweep가 같은 batch를 반복하지 않고 종료하거나 다음 batch로 진행하는 테스트를 추가하세요.

As per path instructions, 비동기 처리 안정성 및 실패 복구 검증 우선 지침을 적용했습니다.

📍 Affects 2 files
  • src/main/java/com/jobdri/jobdri_api/domain/analysis/infrastructure/async/AnalysisAsyncTaskSweepCoordinator.java#L80-L103 (this comment)
  • src/test/java/com/jobdri/jobdri_api/domain/analysis/infrastructure/async/AnalysisAsyncTaskSweepCoordinatorTest.java#L94-L97
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In
`@src/main/java/com/jobdri/jobdri_api/domain/analysis/infrastructure/async/AnalysisAsyncTaskSweepCoordinator.java`
around lines 80 - 103, The sweepTimedOutTaskIds loop can repeatedly load the
same failed first batch because sweepTimedOutTask leaves failed tasks eligible
for selection. Update the coordinator to use keyset pagination or track and
exclude task IDs already processed during the current sweep, ensuring progress
and termination; add a test in AnalysisAsyncTaskSweepCoordinatorTest covering
all 100 IDs failing and verifying the sweep does not repeat the batch and either
terminates or advances to the next batch.

Source: Path instructions


private int sweepTimedOutTaskInTransaction(String taskId) {
AnalysisAsyncTask task = analysisAsyncTaskRepository.findByIdForUpdate(taskId).orElse(null);
if (task == null || task.getStatus() == AnalysisAsyncTaskStatus.SUCCEEDED || task.getStatus() == AnalysisAsyncTaskStatus.FAILED) {
return 0;
Expand All @@ -66,7 +107,7 @@ private int sweepTimedOutTask(String taskId) {
return 0;
}

releaseCreditIfNeeded(task);
analysisAsyncCreditCoordinator.releaseReservedCreditIfNeeded(task);
analysisAsyncTaskService.markFailed(
task.getTaskId(),
expirationDecision.failureReason(),
Expand All @@ -77,7 +118,7 @@ private int sweepTimedOutTask(String taskId) {
}

private ExpirationDecision resolveExpiration(AnalysisAsyncTask task) {
LocalDateTime now = LocalDateTime.now();
LocalDateTime now = LocalDateTime.now(clock);
if (task.getStatus() == AnalysisAsyncTaskStatus.PENDING
&& isExpired(task.getSubmittedAt(), now, analysisQueueProperties.getQueueTimeoutSeconds())) {
return new ExpirationDecision(
Expand Down Expand Up @@ -105,16 +146,11 @@ private boolean isExpired(LocalDateTime baseTime, LocalDateTime now, long timeou
return Duration.between(baseTime, now).getSeconds() >= timeoutSeconds;
}

private void releaseCreditIfNeeded(AnalysisAsyncTask task) {
if (task.getCreditStatus() != AnalysisAsyncCreditStatus.RESERVED || task.getCreditReferenceId() == null) {
return;
}

User user = userService.getUser(task.getUserId());
analysisCreditService.refund(user, task.getCreditReferenceId());
analysisAsyncTaskService.markCreditReleased(task.getTaskId());
private record ExpirationDecision(AnalysisAsyncFailureReason failureReason, String errorMessage) {
}

private record ExpirationDecision(AnalysisAsyncFailureReason failureReason, String errorMessage) {
@FunctionalInterface
private interface TaskIdBatchLoader {
List<String> load();
}
}
Loading
Loading