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,21 +1,42 @@
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,
AnalysisAsyncProcessor analysisAsyncProcessor,
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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -18,7 +19,7 @@ public AnalysisAsyncSweepService(
AnalysisQueueProperties analysisQueueProperties,
Clock clock
) {
super(
this.analysisAsyncTaskSweepCoordinator = new AnalysisAsyncTaskSweepCoordinator(
analysisAsyncTaskRepository,
analysisAsyncTaskService,
analysisAsyncCreditCoordinator,
Expand All @@ -27,4 +28,8 @@ public AnalysisAsyncSweepService(
clock
);
}

public int sweepTimedOutTasks() {
return analysisAsyncTaskSweepCoordinator.sweepTimedOutTasks();
}
}
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -26,7 +35,7 @@ public AnalysisWorkerBridgeService(
ObjectMapper objectMapper,
TransactionTemplate transactionTemplate
) {
super(
this.analysisAsyncWorkerBridge = new AnalysisAsyncWorkerBridge(
analysisAsyncTaskService,
analysisAsyncTaskRepository,
analysisService,
Expand All @@ -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);
}
}
Loading