From 51feca36d07eda76dede659783c01db47fdb0a42 Mon Sep 17 00:00:00 2001 From: backdev-1 Date: Wed, 13 May 2026 16:00:35 +0700 Subject: [PATCH] dev add context for notifcations --- .../entity/search/GlobalSearchTask.java | 3 + .../exception/TaskCancelledException.java | 7 ++ .../search/GlobalSearchAsyncProcessor.java | 107 ++++++++++++++---- .../service/search/GlobalSearchService.java | 2 +- 4 files changed, 97 insertions(+), 22 deletions(-) create mode 100644 src/main/java/ru/soune/nocopy/exception/TaskCancelledException.java diff --git a/src/main/java/ru/soune/nocopy/entity/search/GlobalSearchTask.java b/src/main/java/ru/soune/nocopy/entity/search/GlobalSearchTask.java index 2a7b4c7..e297700 100644 --- a/src/main/java/ru/soune/nocopy/entity/search/GlobalSearchTask.java +++ b/src/main/java/ru/soune/nocopy/entity/search/GlobalSearchTask.java @@ -31,4 +31,7 @@ public class GlobalSearchTask { private LocalDateTime createdAt; private LocalDateTime updatedAt; + + @Version + private Long version; } diff --git a/src/main/java/ru/soune/nocopy/exception/TaskCancelledException.java b/src/main/java/ru/soune/nocopy/exception/TaskCancelledException.java new file mode 100644 index 0000000..9b0295c --- /dev/null +++ b/src/main/java/ru/soune/nocopy/exception/TaskCancelledException.java @@ -0,0 +1,7 @@ +package ru.soune.nocopy.exception; + +public class TaskCancelledException extends RuntimeException { + public TaskCancelledException(String message) { + super(message); + } +} diff --git a/src/main/java/ru/soune/nocopy/service/search/GlobalSearchAsyncProcessor.java b/src/main/java/ru/soune/nocopy/service/search/GlobalSearchAsyncProcessor.java index 0382188..030ddd1 100644 --- a/src/main/java/ru/soune/nocopy/service/search/GlobalSearchAsyncProcessor.java +++ b/src/main/java/ru/soune/nocopy/service/search/GlobalSearchAsyncProcessor.java @@ -2,6 +2,7 @@ package ru.soune.nocopy.service.search; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; +import jakarta.mail.MessagingException; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.scheduling.annotation.Async; @@ -16,6 +17,7 @@ import ru.soune.nocopy.entity.search.GlobalSearchResult; import ru.soune.nocopy.entity.search.GlobalSearchTask; import ru.soune.nocopy.entity.search.SearchStatus; import ru.soune.nocopy.entity.tokenoperation.OperationType; +import ru.soune.nocopy.exception.TaskCancelledException; import ru.soune.nocopy.repository.FileEntityRepository; import ru.soune.nocopy.repository.GlobalSearchResultRepository; import ru.soune.nocopy.repository.GlobalSearchTaskRepository; @@ -53,35 +55,98 @@ public class GlobalSearchAsyncProcessor { private final FileEntityRepository fileEntityRepository; +// @Async +// public void processFilesAsync(String taskId, List filesToProcess, Long userId) { +// try { +// for (int i = 0; i < filesToProcess.size(); i++) { +// String uuid = filesToProcess.get(i); +// log.info("UUID: " + uuid); +// FileEntity file = fileEntityRepository.findById(uuid).orElse(null); +// +// if (file == null) continue; +// +// tariffInfoService.writeOffTokens(userId, TariffConstants.TOKEN_VALUE_FOR_SEARCH, +// OperationType.GLOBAL_SEARCH); +// GlobalSearchResult result = processFile(file, taskId, userId); +// globalSearchResultRepository.save(result); +// +// task.setProcessedFiles(i + 1); +// task.setUpdatedAt(LocalDateTime.now()); +// globalSearchTaskRepository.save(task); +// } +// +// task.setStatus(SearchStatus.COMPLETED.name()); +// task.setUpdatedAt(LocalDateTime.now()); +// GlobalSearchTask save = globalSearchTaskRepository.save(task); +// +// addSearchResultNotification(save, filesToProcess, userId); +// } catch (Exception e) { +// task.setStatus(SearchStatus.FAILED.name()); +// task.setUpdatedAt(LocalDateTime.now()); +// globalSearchTaskRepository.save(task); +// } +// } + @Async - public void processFilesAsync(GlobalSearchTask task, List filesToProcess, Long userId) { + public void processFilesAsync(String taskId, List filesToProcess, Long userId) { try { for (int i = 0; i < filesToProcess.size(); i++) { - String uuid = filesToProcess.get(i); - log.info("UUID: " + uuid); - FileEntity file = fileEntityRepository.findById(uuid).orElse(null); - - if (file == null) continue; - - tariffInfoService.writeOffTokens(userId, TariffConstants.TOKEN_VALUE_FOR_SEARCH, - OperationType.GLOBAL_SEARCH); - GlobalSearchResult result = processFile(file, task.getTaskId(), userId); - globalSearchResultRepository.save(result); - - task.setProcessedFiles(i + 1); - task.setUpdatedAt(LocalDateTime.now()); - globalSearchTaskRepository.save(task); + processOneFile(taskId, filesToProcess.get(i), userId, i); } - - task.setStatus(SearchStatus.COMPLETED.name()); - task.setUpdatedAt(LocalDateTime.now()); - GlobalSearchTask save = globalSearchTaskRepository.save(task); - - addSearchResultNotification(save, filesToProcess, userId); + completeTask(taskId, filesToProcess, userId); + } catch (TaskCancelledException e) { + log.info("Task {} was cancelled", taskId); } catch (Exception e) { + log.error("Task {} failed", taskId, e); + failTask(taskId); + } + } + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void processOneFile(String taskId, String fileId, Long userId, int index) throws MessagingException, IOException { + GlobalSearchTask task = globalSearchTaskRepository.findById(taskId).orElseThrow(); + + if (SearchStatus.FAILED.name().equals(task.getStatus())) { + throw new TaskCancelledException("Task " + taskId + " was cancelled"); + } + + FileEntity file = fileEntityRepository.findById(fileId).orElse(null); + if (file == null) { + log.warn("File not found: {}", fileId); + return; + } + + tariffInfoService.writeOffTokens(userId, TariffConstants.TOKEN_VALUE_FOR_SEARCH, OperationType.GLOBAL_SEARCH); + + GlobalSearchResult result = processFile(file, taskId, userId); + globalSearchResultRepository.save(result); + + task.setProcessedFiles(index + 1); + task.setUpdatedAt(LocalDateTime.now()); + globalSearchTaskRepository.save(task); + } + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void completeTask(String taskId, List filesToProcess, Long userId) { + GlobalSearchTask task = globalSearchTaskRepository.findById(taskId).orElseThrow(); + task.setStatus(SearchStatus.COMPLETED.name()); + task.setUpdatedAt(LocalDateTime.now()); + + GlobalSearchTask saved = globalSearchTaskRepository.save(task); + + addSearchResultNotification(saved, filesToProcess, userId); + } + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void failTask(String taskId) { + try { + GlobalSearchTask task = globalSearchTaskRepository.findById(taskId).orElseThrow(); task.setStatus(SearchStatus.FAILED.name()); task.setUpdatedAt(LocalDateTime.now()); + globalSearchTaskRepository.save(task); + } catch (Exception ex) { + log.error("Could not mark task {} as failed", taskId, ex); } } diff --git a/src/main/java/ru/soune/nocopy/service/search/GlobalSearchService.java b/src/main/java/ru/soune/nocopy/service/search/GlobalSearchService.java index 37d2528..1d7095f 100644 --- a/src/main/java/ru/soune/nocopy/service/search/GlobalSearchService.java +++ b/src/main/java/ru/soune/nocopy/service/search/GlobalSearchService.java @@ -48,7 +48,7 @@ public class GlobalSearchService { List fileIds = filesToProcess.stream().map(FileEntity::getId).toList(); - asyncProcessor.processFilesAsync(save, fileIds, userId); + asyncProcessor.processFilesAsync(save.getTaskId(), fileIds, userId); return save.getTaskId(); }