diff --git a/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncFacadeService.java b/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncFacadeService.java index c5ffe75..0d2b5c9 100644 --- a/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncFacadeService.java +++ b/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncFacadeService.java @@ -1,14 +1,18 @@ package com.jobdri.jobdri_api.domain.analysis.service.async; import com.jobdri.jobdri_api.domain.analysis.application.usecase.async.AnalysisAsyncUseCase; +import com.jobdri.jobdri_api.domain.analysis.dto.response.AnalysisAsyncCancelResponse; +import com.jobdri.jobdri_api.domain.analysis.dto.response.AnalysisAsyncStatusResponse; +import com.jobdri.jobdri_api.domain.analysis.dto.response.AnalysisAsyncSubmitResponse; 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; -import org.springframework.dao.DataIntegrityViolationException; import org.springframework.stereotype.Service; @Service // 분석 비동기 작업의 접수와 상태 조회를 외부 API 관점에서 조율하는 서비스다. -public class AnalysisAsyncFacadeService extends AnalysisAsyncUseCase { +public class AnalysisAsyncFacadeService { + private final AnalysisAsyncUseCase analysisAsyncUseCase; public AnalysisAsyncFacadeService( AnalysisAsyncTaskService analysisAsyncTaskService, @@ -16,6 +20,23 @@ public AnalysisAsyncFacadeService( AnalysisService analysisService, UserService userService ) { - super(analysisAsyncTaskService, analysisAsyncProcessor, analysisService, userService); + this.analysisAsyncUseCase = new AnalysisAsyncUseCase( + analysisAsyncTaskService, + analysisAsyncProcessor, + analysisService, + userService + ); + } + + public AnalysisAsyncSubmitResponse submit(User user, Long mockApplyId) { + return analysisAsyncUseCase.submit(user, mockApplyId); + } + + public AnalysisAsyncStatusResponse getTask(User user, Long mockApplyId, String taskId) { + return analysisAsyncUseCase.getTask(user, mockApplyId, taskId); + } + + public AnalysisAsyncCancelResponse cancel(User user, Long mockApplyId, String taskId) { + return analysisAsyncUseCase.cancel(user, mockApplyId, taskId); } } diff --git a/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncProcessor.java b/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncProcessor.java index 2c44076..068a958 100644 --- a/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncProcessor.java +++ b/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncProcessor.java @@ -5,9 +5,14 @@ @Service // 분석 비동기 작업을 MQ 메시지로 변환해 워커 실행 경로로 넘기는 서비스다. -public class AnalysisAsyncProcessor extends AnalysisAsyncQueueProcessor { +public class AnalysisAsyncProcessor { + private final AnalysisAsyncQueueProcessor analysisAsyncQueueProcessor; public AnalysisAsyncProcessor(AnalysisTaskMessagePublisher analysisTaskMessagePublisher) { - super(analysisTaskMessagePublisher); + this.analysisAsyncQueueProcessor = new AnalysisAsyncQueueProcessor(analysisTaskMessagePublisher); + } + + public void process(String taskId, Long userId, Long mockApplyId, int maxRetryCount) { + analysisAsyncQueueProcessor.process(taskId, userId, mockApplyId, maxRetryCount); } } diff --git a/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncSweepService.java b/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncSweepService.java index abee116..32b7726 100644 --- a/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncSweepService.java +++ b/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisAsyncSweepService.java @@ -8,7 +8,8 @@ import java.time.Clock; @Service -public class AnalysisAsyncSweepService extends AnalysisAsyncTaskSweepCoordinator { +public class AnalysisAsyncSweepService { + private final AnalysisAsyncTaskSweepCoordinator analysisAsyncTaskSweepCoordinator; public AnalysisAsyncSweepService( AnalysisAsyncTaskRepository analysisAsyncTaskRepository, @@ -18,7 +19,7 @@ public AnalysisAsyncSweepService( AnalysisQueueProperties analysisQueueProperties, Clock clock ) { - super( + this.analysisAsyncTaskSweepCoordinator = new AnalysisAsyncTaskSweepCoordinator( analysisAsyncTaskRepository, analysisAsyncTaskService, analysisAsyncCreditCoordinator, @@ -27,4 +28,8 @@ public AnalysisAsyncSweepService( clock ); } + + public int sweepTimedOutTasks() { + return analysisAsyncTaskSweepCoordinator.sweepTimedOutTasks(); + } } diff --git a/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisWorkerBridgeService.java b/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisWorkerBridgeService.java index 89c55cd..0ab933f 100644 --- a/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisWorkerBridgeService.java +++ b/src/main/java/com/jobdri/jobdri_api/domain/analysis/service/async/AnalysisWorkerBridgeService.java @@ -1,19 +1,28 @@ package com.jobdri.jobdri_api.domain.analysis.service.async; import com.fasterxml.jackson.databind.ObjectMapper; +import com.jobdri.jobdri_api.domain.analysis.dto.internal.worker.AnalysisWorkerCompleteRequest; +import com.jobdri.jobdri_api.domain.analysis.dto.internal.worker.AnalysisWorkerContextResponse; +import com.jobdri.jobdri_api.domain.analysis.dto.internal.worker.AnalysisWorkerResultStoreRequest; +import com.jobdri.jobdri_api.domain.analysis.dto.response.AnalysisResponse; import com.jobdri.jobdri_api.domain.analysis.infrastructure.async.AnalysisAsyncWorkerBridge; import com.jobdri.jobdri_api.domain.analysis.repository.AnalysisAsyncTaskRepository; +import com.jobdri.jobdri_api.domain.analysis.type.AnalysisAsyncFailureReason; 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.service.UserService; +import com.jobdri.jobdri_api.domain.workerresult.dto.WorkerTaskResultResponse; import com.jobdri.jobdri_api.domain.workerresult.service.WorkerTaskResultService; -import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.support.TransactionTemplate; +import java.time.Instant; + @Service // 외부 분석 워커와 내부 분석 도메인 상태를 연결해 주는 브리지 서비스다. -public class AnalysisWorkerBridgeService extends AnalysisAsyncWorkerBridge { +public class AnalysisWorkerBridgeService { + private final AnalysisAsyncWorkerBridge analysisAsyncWorkerBridge; public AnalysisWorkerBridgeService( AnalysisAsyncTaskService analysisAsyncTaskService, @@ -26,7 +35,7 @@ public AnalysisWorkerBridgeService( ObjectMapper objectMapper, TransactionTemplate transactionTemplate ) { - super( + this.analysisAsyncWorkerBridge = new AnalysisAsyncWorkerBridge( analysisAsyncTaskService, analysisAsyncTaskRepository, analysisService, @@ -38,4 +47,66 @@ public AnalysisWorkerBridgeService( transactionTemplate ); } + + @Transactional + public void markRunning(String taskId, String workerId, int retryCount, Instant submittedAt) { + analysisAsyncWorkerBridge.markRunning(taskId, workerId, retryCount, submittedAt); + } + + @Transactional + public void markRetry( + String taskId, + AnalysisAsyncFailureReason failureReason, + String errorMessage, + int retryCount, + String workerId, + Long queueLatencyMillis + ) { + analysisAsyncWorkerBridge.markRetry( + taskId, + failureReason, + errorMessage, + retryCount, + workerId, + queueLatencyMillis + ); + } + + @Transactional + public void failTask( + String taskId, + AnalysisAsyncFailureReason failureReason, + String errorMessage, + int retryCount, + String workerId, + Long queueLatencyMillis + ) { + analysisAsyncWorkerBridge.failTask( + taskId, + failureReason, + errorMessage, + retryCount, + workerId, + queueLatencyMillis + ); + } + + public AnalysisWorkerContextResponse getContext(String taskId, Long userId, Long mockApplyId) { + return analysisAsyncWorkerBridge.getContext(taskId, userId, mockApplyId); + } + + @Transactional + public AnalysisResponse completeTask(String taskId, AnalysisWorkerCompleteRequest request) { + return analysisAsyncWorkerBridge.completeTask(taskId, request); + } + + @Transactional + public void storeGeneratedResult(String taskId, AnalysisWorkerResultStoreRequest request) { + analysisAsyncWorkerBridge.storeGeneratedResult(taskId, request); + } + + @Transactional(readOnly = true) + public WorkerTaskResultResponse getStoredResult(String taskId) { + return analysisAsyncWorkerBridge.getStoredResult(taskId); + } }