From 31c5d9aaabbc3fba88e1e29b7d0713e4fad723bb Mon Sep 17 00:00:00 2001 From: vladp Date: Fri, 6 Mar 2026 17:28:52 +0700 Subject: [PATCH] dev async global search --- .../configuration/file/AsyncConfig.java | 11 ++ .../controller/GlobalSearchController.java | 140 +++++++++++++++++ .../dto/search/GlobalSearchFileResult.java | 19 +++ .../dto/search/GlobalSearchStartRequest.java | 14 ++ .../dto/search/GlobalSearchStartResponse.java | 12 ++ .../search/GlobalSearchStatusResponse.java | 18 +++ .../entity/search/GlobalSearchResult.java | 32 ++++ .../entity/search/GlobalSearchTask.java | 34 ++++ .../nocopy/entity/search/SearchStatus.java | 5 + .../handler/ImageFoundRequestHandler.java | 2 - .../GlobalSearchResultRepository.java | 13 ++ .../GlobalSearchTaskRepository.java | 13 ++ .../search/GlobalSearchAsyncProcessor.java | 147 ++++++++++++++++++ .../service/search/GlobalSearchService.java | 62 ++++++++ 14 files changed, 520 insertions(+), 2 deletions(-) create mode 100644 src/main/java/ru/soune/nocopy/controller/GlobalSearchController.java create mode 100644 src/main/java/ru/soune/nocopy/dto/search/GlobalSearchFileResult.java create mode 100644 src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStartRequest.java create mode 100644 src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStartResponse.java create mode 100644 src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStatusResponse.java create mode 100644 src/main/java/ru/soune/nocopy/entity/search/GlobalSearchResult.java create mode 100644 src/main/java/ru/soune/nocopy/entity/search/GlobalSearchTask.java create mode 100644 src/main/java/ru/soune/nocopy/entity/search/SearchStatus.java create mode 100644 src/main/java/ru/soune/nocopy/repository/GlobalSearchResultRepository.java create mode 100644 src/main/java/ru/soune/nocopy/repository/GlobalSearchTaskRepository.java create mode 100644 src/main/java/ru/soune/nocopy/service/search/GlobalSearchAsyncProcessor.java create mode 100644 src/main/java/ru/soune/nocopy/service/search/GlobalSearchService.java diff --git a/src/main/java/ru/soune/nocopy/configuration/file/AsyncConfig.java b/src/main/java/ru/soune/nocopy/configuration/file/AsyncConfig.java index 6882014..d90475e 100644 --- a/src/main/java/ru/soune/nocopy/configuration/file/AsyncConfig.java +++ b/src/main/java/ru/soune/nocopy/configuration/file/AsyncConfig.java @@ -21,4 +21,15 @@ public class AsyncConfig { executor.initialize(); return executor; } + + @Bean(name = "taskExecutor") + public Executor taskExecutor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(3); + executor.setMaxPoolSize(5); + executor.setQueueCapacity(50); + executor.setThreadNamePrefix("global-search-"); + executor.initialize(); + return executor; + } } diff --git a/src/main/java/ru/soune/nocopy/controller/GlobalSearchController.java b/src/main/java/ru/soune/nocopy/controller/GlobalSearchController.java new file mode 100644 index 0000000..cab9462 --- /dev/null +++ b/src/main/java/ru/soune/nocopy/controller/GlobalSearchController.java @@ -0,0 +1,140 @@ +package ru.soune.nocopy.controller; + +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.*; +import ru.soune.nocopy.dto.file.YandexSearchResponse; +import ru.soune.nocopy.dto.search.GlobalSearchFileResult; +import ru.soune.nocopy.dto.search.GlobalSearchStartRequest; +import ru.soune.nocopy.dto.search.GlobalSearchStartResponse; +import ru.soune.nocopy.dto.search.GlobalSearchStatusResponse; +import ru.soune.nocopy.entity.file.FileEntity; +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.repository.GlobalSearchResultRepository; +import ru.soune.nocopy.repository.GlobalSearchTaskRepository; +import ru.soune.nocopy.service.register.AuthService; +import ru.soune.nocopy.service.search.GlobalSearchService; + +import java.util.*; +import java.util.stream.Collectors; + +@Slf4j +@RestController +@RequestMapping("/api/v1/global-search") +@RequiredArgsConstructor +public class GlobalSearchController { + + private final AuthService authService; + + private final GlobalSearchService globalSearchService; + + private final GlobalSearchTaskRepository globalSearchTaskRepository; + + private final GlobalSearchTaskRepository searchTaskRepository; + + private final GlobalSearchResultRepository globalSearchResultRepository; + + private final ObjectMapper objectMapper; + + @PostMapping("/start") + public ResponseEntity startSearch( + @RequestBody GlobalSearchStartRequest request, + @RequestHeader(value = "Authorization", required = false) String tokenHeader) { + + if (tokenHeader == null || tokenHeader.isBlank()) { + Map errorData = new HashMap<>(); + errorData.put("token", tokenHeader); + + return ResponseEntity.ok().body(Map.of("error", errorData)); + } + + Long userId = authService.useUserAuthToken(tokenHeader); + + List filesToProcess = globalSearchService.getFilesToProcess(request, userId); + + if (filesToProcess.isEmpty()) return ResponseEntity.ok().body("Files for search not found"); + + String taskId = globalSearchService.startSearch(request, userId, filesToProcess); + + Optional taskOptional = globalSearchTaskRepository.findById(taskId); + + if (taskOptional.isEmpty()) return ResponseEntity.ok().body("Task not found"); + + GlobalSearchTask task = taskOptional.orElseThrow(); + GlobalSearchStartResponse response = new GlobalSearchStartResponse(); + response.setTaskId(taskId); + response.setStatus(SearchStatus.ACCEPTED.name()); + response.setTotalFiles(task.getTotalFiles()); + + return ResponseEntity.ok(response); + } + + @GetMapping("/status/{taskId}") + public ResponseEntity getStatus( + @PathVariable String taskId, + @RequestHeader(value = "Authorization", required = false) String tokenHeader) { + + if (tokenHeader == null || tokenHeader.isBlank()) { + Map errorData = new HashMap<>(); + errorData.put("token", tokenHeader); + + return ResponseEntity.ok().body(Map.of("error", errorData)); + } + + Long userId = authService.useUserAuthToken(tokenHeader); + + Optional taskOptional = searchTaskRepository.findByTaskIdAndUserId(taskId, userId); + + if (taskOptional.isEmpty()) return ResponseEntity.ok().body("Task not found"); + + GlobalSearchTask task = taskOptional.orElseThrow(); + + List results = globalSearchResultRepository.findByTaskIdOrderByIdAsc(taskId); + + List fileResults = results.stream() + .map(result -> { + GlobalSearchFileResult dto = new GlobalSearchFileResult(); + dto.setFileId(result.getFileId()); + dto.setFileName(result.getFileName()); + dto.setThumbnail(result.getThumbnail()); + dto.setFileStatus(result.getFileStatus()); + + try { + YandexSearchResponse searchResponse = objectMapper.readValue( + result.getSearchResults(), + YandexSearchResponse.class + ); + + dto.setImages(searchResponse.getImages()); + dto.setPage(searchResponse.getPage()); + dto.setPageSize(searchResponse.getPageSize()); + dto.setTotalResults(searchResponse.getTotalResults()); + dto.setTotalPages(searchResponse.getTotalPages()); + + } catch (Exception e) { + log.error("Failed to parse search results for file: {}", result.getFileId(), e); + dto.setImages(new ArrayList<>()); + dto.setTotalResults(0); + } + + return dto; + }) + .collect(Collectors.toList()); + + int progress = task.getTotalFiles() > 0 + ? task.getProcessedFiles() * 100 / task.getTotalFiles() + : 0; + + GlobalSearchStatusResponse response = new GlobalSearchStatusResponse(); + response.setTaskId(taskId); + response.setStatus(task.getStatus()); + response.setProgress(progress); + response.setResults(fileResults); + + return ResponseEntity.ok(response); + } +} diff --git a/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchFileResult.java b/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchFileResult.java new file mode 100644 index 0000000..31d1358 --- /dev/null +++ b/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchFileResult.java @@ -0,0 +1,19 @@ +package ru.soune.nocopy.dto.search; + +import lombok.Data; +import ru.soune.nocopy.dto.file.YandexSearchResponse; + +import java.util.List; + +@Data +public class GlobalSearchFileResult { + private String fileId; + private String fileName; + private String thumbnail; + private String fileStatus; + private List images; + private int page; + private int pageSize; + private int totalResults; + private int totalPages; +} diff --git a/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStartRequest.java b/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStartRequest.java new file mode 100644 index 0000000..9e1021a --- /dev/null +++ b/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStartRequest.java @@ -0,0 +1,14 @@ +package ru.soune.nocopy.dto.search; + +import lombok.Builder; +import lombok.Data; + +import java.util.List; + +@Builder +@Data +public class GlobalSearchStartRequest { + private String searchType; + + private List fileIds; +} diff --git a/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStartResponse.java b/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStartResponse.java new file mode 100644 index 0000000..9b55562 --- /dev/null +++ b/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStartResponse.java @@ -0,0 +1,12 @@ +package ru.soune.nocopy.dto.search; + +import lombok.Data; + +@Data +public class GlobalSearchStartResponse { + private String taskId; + + private String status; + + private Integer totalFiles; +} diff --git a/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStatusResponse.java b/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStatusResponse.java new file mode 100644 index 0000000..0634aa1 --- /dev/null +++ b/src/main/java/ru/soune/nocopy/dto/search/GlobalSearchStatusResponse.java @@ -0,0 +1,18 @@ +package ru.soune.nocopy.dto.search; + +import lombok.Data; + +import java.util.List; + +@Data +public class GlobalSearchStatusResponse { + private String taskId; + + private String status; + + private String searchType; + + private Integer progress; + + private List results; +} diff --git a/src/main/java/ru/soune/nocopy/entity/search/GlobalSearchResult.java b/src/main/java/ru/soune/nocopy/entity/search/GlobalSearchResult.java new file mode 100644 index 0000000..ab8661f --- /dev/null +++ b/src/main/java/ru/soune/nocopy/entity/search/GlobalSearchResult.java @@ -0,0 +1,32 @@ +package ru.soune.nocopy.entity.search; + +import jakarta.persistence.*; +import lombok.Data; +import lombok.Getter; +import lombok.RequiredArgsConstructor; +import lombok.Setter; + +@Entity +@Table(name = "global_search_results") +@Data +@Getter @Setter +@RequiredArgsConstructor +public class GlobalSearchResult { + + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + private String taskId; + + private String fileId; + + private String fileName; + + private String thumbnail; + + private String fileStatus; + + @Column(length = 5000) + private String searchResults; +} diff --git a/src/main/java/ru/soune/nocopy/entity/search/GlobalSearchTask.java b/src/main/java/ru/soune/nocopy/entity/search/GlobalSearchTask.java new file mode 100644 index 0000000..2a7b4c7 --- /dev/null +++ b/src/main/java/ru/soune/nocopy/entity/search/GlobalSearchTask.java @@ -0,0 +1,34 @@ +package ru.soune.nocopy.entity.search; + +import jakarta.persistence.*; +import lombok.Data; +import lombok.Getter; +import lombok.RequiredArgsConstructor; +import lombok.Setter; + +import java.time.LocalDateTime; + +@Entity +@Table(name = "global_search_tasks") +@Data +@Getter @Setter +@RequiredArgsConstructor +public class GlobalSearchTask { + @Id + @GeneratedValue(strategy = GenerationType.UUID) + private String taskId; + + private Long userId; + + private String searchType; + + private String status; + + private Integer totalFiles; + + private Integer processedFiles; + + private LocalDateTime createdAt; + + private LocalDateTime updatedAt; +} diff --git a/src/main/java/ru/soune/nocopy/entity/search/SearchStatus.java b/src/main/java/ru/soune/nocopy/entity/search/SearchStatus.java new file mode 100644 index 0000000..867abb8 --- /dev/null +++ b/src/main/java/ru/soune/nocopy/entity/search/SearchStatus.java @@ -0,0 +1,5 @@ +package ru.soune.nocopy.entity.search; + +public enum SearchStatus { + PROCESSING, COMPLETED, ACCEPTED, FAILED, SUCCESS, TIMEOUT +} diff --git a/src/main/java/ru/soune/nocopy/handler/ImageFoundRequestHandler.java b/src/main/java/ru/soune/nocopy/handler/ImageFoundRequestHandler.java index 81182da..01e8a27 100644 --- a/src/main/java/ru/soune/nocopy/handler/ImageFoundRequestHandler.java +++ b/src/main/java/ru/soune/nocopy/handler/ImageFoundRequestHandler.java @@ -90,10 +90,8 @@ public class ImageFoundRequestHandler implements RequestHandler { allUniqueImages = searchImageService.removeDuplicateUrls(allYandexImages, allGoogleImages); log.info("allUniqueImages OK: {} images", allUniqueImages.size()); -// if (!allUniqueImages.isEmpty()) { tariffInfoService.writeOffTokens(fileEntity.getUserId(), TariffConstants.TOKEN_VALUE_FOR_SEARCH); checkCounterService.incrementCheckCount(fileEntity.getUserId(), fileEntity.getMimeType()); -// } int page = imageSearchRequest.getPage() != null ? imageSearchRequest.getPage() : 1; int pageSize = 5; diff --git a/src/main/java/ru/soune/nocopy/repository/GlobalSearchResultRepository.java b/src/main/java/ru/soune/nocopy/repository/GlobalSearchResultRepository.java new file mode 100644 index 0000000..dc41337 --- /dev/null +++ b/src/main/java/ru/soune/nocopy/repository/GlobalSearchResultRepository.java @@ -0,0 +1,13 @@ +package ru.soune.nocopy.repository; + +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.stereotype.Repository; +import ru.soune.nocopy.entity.search.GlobalSearchResult; + +import java.util.List; + +@Repository +public interface GlobalSearchResultRepository extends JpaRepository { + List findByTaskIdOrderByIdAsc(String taskId); + void deleteByTaskId(String taskId); +} diff --git a/src/main/java/ru/soune/nocopy/repository/GlobalSearchTaskRepository.java b/src/main/java/ru/soune/nocopy/repository/GlobalSearchTaskRepository.java new file mode 100644 index 0000000..908a659 --- /dev/null +++ b/src/main/java/ru/soune/nocopy/repository/GlobalSearchTaskRepository.java @@ -0,0 +1,13 @@ +package ru.soune.nocopy.repository; + +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.stereotype.Repository; +import ru.soune.nocopy.entity.search.GlobalSearchTask; + +import java.util.Optional; + +@Repository +public interface GlobalSearchTaskRepository extends JpaRepository { + Optional findByTaskIdAndUserId(String taskId, Long userId); +} + diff --git a/src/main/java/ru/soune/nocopy/service/search/GlobalSearchAsyncProcessor.java b/src/main/java/ru/soune/nocopy/service/search/GlobalSearchAsyncProcessor.java new file mode 100644 index 0000000..cd55272 --- /dev/null +++ b/src/main/java/ru/soune/nocopy/service/search/GlobalSearchAsyncProcessor.java @@ -0,0 +1,147 @@ +package ru.soune.nocopy.service.search; + +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.scheduling.annotation.Async; +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Transactional; +import ru.soune.nocopy.dto.file.YandexSearchResponse; +import ru.soune.nocopy.entity.file.FileEntity; +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.repository.GlobalSearchResultRepository; +import ru.soune.nocopy.repository.GlobalSearchTaskRepository; +import ru.soune.nocopy.service.tariff.TariffConstants; +import ru.soune.nocopy.service.tariff.TariffInfoService; + +import java.io.IOException; +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.TimeoutException; +import java.util.stream.Collectors; + +@Slf4j +@Component +@RequiredArgsConstructor +public class GlobalSearchAsyncProcessor { + + private final GlobalSearchResultRepository globalSearchResultRepository; + + private final SearchImageService searchImageService; + + private final ObjectMapper objectMapper; + + private final TariffInfoService tariffInfoService; + + private final GlobalSearchTaskRepository globalSearchTaskRepository; + + @Async + @Transactional + public void processFilesAsync(String taskId, List filesToProcess, Long userId) { + try { + for (int i = 0; i < filesToProcess.size(); i++) { + FileEntity file = filesToProcess.get(i); + + GlobalSearchResult result = processFile(file, taskId, userId); + globalSearchResultRepository.save(result); + + GlobalSearchTask task = globalSearchTaskRepository.findById(taskId).orElseThrow(); + task.setProcessedFiles(i + 1); + task.setUpdatedAt(LocalDateTime.now()); + globalSearchTaskRepository.save(task); + } + + GlobalSearchTask task = globalSearchTaskRepository.findById(taskId).orElseThrow(); + task.setStatus(SearchStatus.COMPLETED.name()); + task.setUpdatedAt(LocalDateTime.now()); + globalSearchTaskRepository.save(task); + + } catch (Exception e) { + log.error("Global search failed for task: {}", taskId, e); + + GlobalSearchTask task = globalSearchTaskRepository.findById(taskId).orElseThrow(); + task.setStatus(SearchStatus.FAILED.name()); + task.setUpdatedAt(LocalDateTime.now()); + globalSearchTaskRepository.save(task); + } + } + + private GlobalSearchResult processFile(FileEntity file, String taskId, Long userId) { + GlobalSearchResult result = new GlobalSearchResult(); + result.setTaskId(taskId); + result.setFileId(file.getId()); + result.setFileName(file.getOriginalFileName()); + result.setThumbnail(file.getThumbnailPath()); + + List allUniqueImages = new ArrayList<>(); + boolean hasTimeout = false; + + try { + String searchResponseGoogle = searchImageService.searchReverseByPublicUrl( + file, "google_lens", "exact_matches"); + + List googleImages = + searchImageService.getAllImagesWithoutPagination(searchResponseGoogle, "exact_matches"); + + allUniqueImages.addAll(googleImages); + log.info("Google search OK for file {}", file.getId()); + + } catch (TimeoutException e) { + log.warn("Google search timeout for file {}", file.getId()); + hasTimeout = true; + } catch (IOException e) { + log.error("Google search failed for file {}", file.getId(), e); + } + + try { + String searchResponseYandex = searchImageService.searchReverseByPublicUrl( + file, "yandex_reverse_image", "visual_matches"); + + List yandexImages = + searchImageService.getAllImagesWithoutPagination(searchResponseYandex, "visual_matches"); + + allUniqueImages.addAll(yandexImages); + log.info("Yandex search OK for file {}", file.getId()); + + } catch (TimeoutException e) { + log.warn("Yandex search timeout for file {}", file.getId()); + hasTimeout = true; + } catch (IOException e) { + log.error("Yandex search failed for file {}", file.getId(), e); + } + + allUniqueImages = searchImageService.removeDuplicateUrls( + allUniqueImages.stream() + .filter(img -> img.getUrl() != null) + .collect(Collectors.toList()), + new ArrayList<>() + ); + + YandexSearchResponse searchResponse = new YandexSearchResponse(); + searchResponse.setImages(allUniqueImages); + searchResponse.setPage(1); + searchResponse.setPageSize(allUniqueImages.size()); + searchResponse.setTotalResults(allUniqueImages.size()); + searchResponse.setTotalPages(1); + + try { + result.setSearchResults(objectMapper.writeValueAsString(searchResponse)); + } catch (Exception e) { + log.error("Failed to serialize search results", e); + result.setSearchResults("{}"); + } + + if (allUniqueImages.isEmpty() && hasTimeout) { + result.setFileStatus(SearchStatus.TIMEOUT.name()); + } else { + result.setFileStatus(SearchStatus.SUCCESS.name()); + + tariffInfoService.writeOffTokens(userId, TariffConstants.TOKEN_VALUE_FOR_SEARCH); + } + + return result; + } +} diff --git a/src/main/java/ru/soune/nocopy/service/search/GlobalSearchService.java b/src/main/java/ru/soune/nocopy/service/search/GlobalSearchService.java new file mode 100644 index 0000000..4b5ae40 --- /dev/null +++ b/src/main/java/ru/soune/nocopy/service/search/GlobalSearchService.java @@ -0,0 +1,62 @@ +package ru.soune.nocopy.service.search; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; +import ru.soune.nocopy.dto.search.GlobalSearchStartRequest; +import ru.soune.nocopy.entity.file.FileEntity; +import ru.soune.nocopy.entity.monitoring.FileMonitoringEntity; +import ru.soune.nocopy.entity.search.GlobalSearchTask; +import ru.soune.nocopy.entity.search.SearchStatus; +import ru.soune.nocopy.repository.FileEntityRepository; +import ru.soune.nocopy.repository.FileMonitoringRepository; +import ru.soune.nocopy.repository.GlobalSearchTaskRepository; + +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.List; + +@Slf4j +@Service +@RequiredArgsConstructor +public class GlobalSearchService { + + private final FileEntityRepository fileEntityRepository; + + private final FileMonitoringRepository fileMonitoringRepository; + + private final GlobalSearchTaskRepository globalSearchTaskRepository; + + private final GlobalSearchAsyncProcessor asyncProcessor; + + @Transactional + public String startSearch(GlobalSearchStartRequest request, Long userId, List filesToProcess) { + GlobalSearchTask task = new GlobalSearchTask(); + task.setUserId(userId); + task.setSearchType(request.getSearchType()); + task.setStatus(SearchStatus.PROCESSING.name()); + task.setTotalFiles(filesToProcess.size()); + task.setProcessedFiles(0); + task.setCreatedAt(LocalDateTime.now()); + + globalSearchTaskRepository.save(task); + + asyncProcessor.processFilesAsync(task.getTaskId(), filesToProcess, userId); + + return task.getTaskId(); + } + + public List getFilesToProcess(GlobalSearchStartRequest request, Long userId) { + if ("monitoring".equals(request.getSearchType())) { + List monitoringFiles = + fileMonitoringRepository.findByUserIdAndIsActiveTrue(userId); + + return monitoringFiles.isEmpty() ? new ArrayList<>(): monitoringFiles.stream() + .map(FileMonitoringEntity::getFile) + .toList(); + } else { + return fileEntityRepository.findAllById(request.getFileIds()); + } + } +}