From 76b678460a202608941b5715cb8ddb2f60faa8ea Mon Sep 17 00:00:00 2001 From: backdev-1 Date: Thu, 11 Jun 2026 11:07:34 +0700 Subject: [PATCH] add work monitoring scheme --- build.gradle | 5 + docker-compose.yaml | 14 +- .../no_copy/monitoring/dto/MonitoringDTO.java | 17 ++ .../monitoring/dto/SearchResponse.java | 47 +++++ .../kafka/MonitoringCommandListener.java | 20 +- .../searcher/SearchImageService.java | 172 ++++++++++++++++++ .../monitoring/service/MonitoringService.java | 45 ++++- src/main/resources/application.yaml | 8 +- 8 files changed, 313 insertions(+), 15 deletions(-) create mode 100644 src/main/java/ru/no_copy/monitoring/dto/MonitoringDTO.java create mode 100644 src/main/java/ru/no_copy/monitoring/dto/SearchResponse.java create mode 100644 src/main/java/ru/no_copy/monitoring/searcher/SearchImageService.java diff --git a/build.gradle b/build.gradle index 39f5a43..066d8fc 100644 --- a/build.gradle +++ b/build.gradle @@ -42,6 +42,11 @@ dependencies { testImplementation 'org.springframework.boot:spring-boot-starter-thymeleaf-test' testImplementation 'org.springframework.boot:spring-boot-starter-validation-test' testRuntimeOnly 'org.junit.platform:junit-platform-launcher' + implementation 'com.fasterxml.jackson.dataformat:jackson-dataformat-xml:2.17.0' + implementation 'tools.jackson.core:jackson-core:3.0.3' + implementation group: 'com.fasterxml.jackson.core', name: 'jackson-databind', version: '2.17.2' + implementation group: 'com.squareup.okhttp3', name: 'okhttp', version: '4.12.0' + } tasks.named('test') { diff --git a/docker-compose.yaml b/docker-compose.yaml index 88f083e..7ec1267 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -27,10 +27,12 @@ services: POSTGRES_USER: adminMonitoring POSTGRES_PASSWORD: monitoringDbApp SERVER_PORT: ${SERVER_PORT:-8083} - SPRING_KAFKA_BOOTSTRAP_SERVERS: kafka:9092 +# SPRING_KAFKA_BOOTSTRAP_SERVERS: kafka:9092 depends_on: - db restart: unless-stopped + networks: + - app-network kafka: image: apache/kafka:latest @@ -52,6 +54,8 @@ services: KAFKA_MIN_INSYNC_REPLICAS: 1 volumes: - kafka_data:/var/lib/kafka/data + networks: + - app-network kafka-ui: image: provectuslabs/kafka-ui:latest @@ -64,8 +68,14 @@ services: KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: 'kafka:9092' depends_on: - kafka + networks: + - app-network volumes: kafka_data: - postgres_data: \ No newline at end of file + postgres_data: + +networks: + app-network: + external: true \ No newline at end of file diff --git a/src/main/java/ru/no_copy/monitoring/dto/MonitoringDTO.java b/src/main/java/ru/no_copy/monitoring/dto/MonitoringDTO.java new file mode 100644 index 0000000..a1c1218 --- /dev/null +++ b/src/main/java/ru/no_copy/monitoring/dto/MonitoringDTO.java @@ -0,0 +1,17 @@ +package ru.no_copy.monitoring.dto; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class MonitoringDTO { + String fileId; + String baseUrl; + String engine; + String searchType; +} diff --git a/src/main/java/ru/no_copy/monitoring/dto/SearchResponse.java b/src/main/java/ru/no_copy/monitoring/dto/SearchResponse.java new file mode 100644 index 0000000..34befff --- /dev/null +++ b/src/main/java/ru/no_copy/monitoring/dto/SearchResponse.java @@ -0,0 +1,47 @@ +package ru.no_copy.monitoring.dto; + +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.Data; + +import java.util.List; + +@Data +@JsonIgnoreProperties(ignoreUnknown = true) +public class SearchResponse { + @JsonProperty("images") + private List images; + + @Data + @JsonIgnoreProperties(ignoreUnknown = true) + public static class ImageResult { + @JsonProperty("url") + private String url; + + @JsonProperty("pageUrl") + private String pageUrl; + + @JsonProperty("pageTitle") + private String pageTitle; + + @JsonProperty("width") + private Integer width; + + @JsonProperty("height") + private Integer height; + + @JsonProperty("host") + private String host; + + @JsonProperty("file_id") + private String fileId; + } + + private int page; + + private int pageSize; + + private int totalResults; + + private int totalPages; +} diff --git a/src/main/java/ru/no_copy/monitoring/kafka/MonitoringCommandListener.java b/src/main/java/ru/no_copy/monitoring/kafka/MonitoringCommandListener.java index 3639460..332df85 100644 --- a/src/main/java/ru/no_copy/monitoring/kafka/MonitoringCommandListener.java +++ b/src/main/java/ru/no_copy/monitoring/kafka/MonitoringCommandListener.java @@ -1,30 +1,32 @@ package ru.no_copy.monitoring.kafka; +import lombok.extern.slf4j.Slf4j; import org.springframework.kafka.annotation.KafkaListener; -import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Component; import ru.no_copy.monitoring.service.MonitoringService; +import java.io.IOException; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeoutException; @Component +@Slf4j public class MonitoringCommandListener { - private final KafkaTemplate kafkaTemplate; - private final MonitoringService monitoringService; - - public MonitoringCommandListener(KafkaTemplate kafkaTemplate, MonitoringService monitoringService) { - this.kafkaTemplate = kafkaTemplate; + public MonitoringCommandListener(MonitoringService monitoringService) { this.monitoringService = monitoringService; } @KafkaListener(topics = "monitoring-commands", groupId = "monitoring-service") public void handleCommand(String commandMessage) { + log.info("message: " + commandMessage); CompletableFuture.runAsync(() -> { - String result = monitoringService.execute(commandMessage); - //TODO - kafkaTemplate.send("monitoring-results", result); + try { + monitoringService.execute(commandMessage); + } catch (IOException | TimeoutException e) { + throw new RuntimeException(e); + } }); } } diff --git a/src/main/java/ru/no_copy/monitoring/searcher/SearchImageService.java b/src/main/java/ru/no_copy/monitoring/searcher/SearchImageService.java new file mode 100644 index 0000000..11efe1d --- /dev/null +++ b/src/main/java/ru/no_copy/monitoring/searcher/SearchImageService.java @@ -0,0 +1,172 @@ +package ru.no_copy.monitoring.searcher; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import okhttp3.*; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Service; +import ru.no_copy.monitoring.dto.SearchResponse; + +import java.io.IOException; +import java.net.SocketTimeoutException; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; + +@Slf4j +@Service +@RequiredArgsConstructor +public class SearchImageService { + + private final ObjectMapper objectMapper; + + @Value("${searchapi.api-key:}") + private String searchApiKey; + + public String searchReverseByPublicUrl(String fileId, String baseUrl, String engine, String searchType) + throws IOException, TimeoutException { + + String publicUrl = String.format("%s/api/files/public/%s", baseUrl, fileId); + + try { + return callReverseImageApiByUrl(publicUrl, engine, searchType); + } catch (SocketTimeoutException e) { + log.error("Yandex search timeout after {}", fileId); + throw new TimeoutException("Search timeout"); + } + } + + private String callReverseImageApiByUrl(String imageUrl, String engine, String searchType) throws IOException { + if (searchApiKey == null || searchApiKey.isBlank()) { + throw new IllegalStateException("SearchAPI key not configured"); + } + + OkHttpClient client = createHttpClient(); + + HttpUrl url = HttpUrl.parse("https://www.searchapi.io/api/v1/search") + .newBuilder() + .addQueryParameter("engine", engine) + .addQueryParameter("api_key", searchApiKey) + .addQueryParameter("url", imageUrl) + .addQueryParameter("search_type", searchType) + .addQueryParameter("t_", String.valueOf(System.currentTimeMillis())) + .build(); + + Request request = new Request.Builder() + .url(url) + .header("Accept", "application/json") + .header("User-Agent", "Mozilla/5.0") + .build(); + + long start = System.currentTimeMillis(); + + try (Response response = client.newCall(request).execute()) { + ResponseBody body = response.body(); + long duration = System.currentTimeMillis() - start; + log.info("SearchAPI response code={}, duration={}ms, engine={}", + response.code(), duration, engine); + + if (!response.isSuccessful()) { + String errorBody = response.body() != null ? response.body().string() : "null"; + throw new IOException("API error " + response.code() + ": " + errorBody); + } + + if (body == null) { + throw new IOException("Empty response body"); + } + + return body.string(); + } + } + + public List getAllImagesWithoutPagination(String searchApiJson, String findType, String fileId) + throws IOException { + JsonNode root = objectMapper.readTree(searchApiJson); + JsonNode matches = root.path(findType); + + List allImages = new ArrayList<>(); + + if (matches.isArray()) { + for (JsonNode match : matches) { + SearchResponse.ImageResult result = mapImageResult(match); + + if ("exact_matches".equals(findType)) { + JsonNode thumbnail = match.path("thumbnail"); + if (!thumbnail.isMissingNode()) { + String thumbnailUrl = thumbnail.asText(); + if (thumbnailUrl.startsWith("data:image")) { + result.setUrl(thumbnailUrl); + } else if (!thumbnailUrl.isBlank()) { + result.setUrl(thumbnailUrl); + } + } + + JsonNode imageNode = match.path("image"); + if (!imageNode.isMissingNode()) { + String directUrl = imageNode.path("link").asText(); + if (directUrl != null && !directUrl.isBlank()) { + result.setUrl(directUrl); + } + } + } + + if (result.getUrl() != null && !result.getUrl().isBlank()) { + allImages.add(result); + } + } + } + + return allImages; + } + + private SearchResponse.ImageResult mapImageResult(JsonNode match) { + + SearchResponse.ImageResult result = + new SearchResponse.ImageResult(); + + JsonNode imageNode = match.path("image"); + if (imageNode.isObject()) { + result.setUrl(imageNode.path("link").asText()); + result.setWidth(imageNode.path("width").asInt(0)); + result.setHeight(imageNode.path("height").asInt(0)); + } + + result.setPageUrl(match.path("link").asText()); + result.setPageTitle(match.path("title").asText()); + + String source = match.path("source").asText(); + result.setHost(extractHostFromSource(source)); + + return result; + } + + private OkHttpClient createHttpClient() { + OkHttpClient.Builder builder = new OkHttpClient.Builder() + .connectTimeout(35, TimeUnit.SECONDS) + .writeTimeout(35, TimeUnit.SECONDS) + .readTimeout(35, TimeUnit.SECONDS) + .callTimeout(35, TimeUnit.SECONDS) + .followRedirects(true) + .followSslRedirects(true) + .retryOnConnectionFailure(true); + + return builder.build(); + } + + private String extractHostFromSource(String source) { + if (source == null || source.isEmpty()) { + return ""; + } + source = source.replaceFirst("^(https?://)?(www\\.)?", ""); + + int slashIndex = source.indexOf('/'); + if (slashIndex > 0) { + return source.substring(0, slashIndex); + } + + return source; + } +} diff --git a/src/main/java/ru/no_copy/monitoring/service/MonitoringService.java b/src/main/java/ru/no_copy/monitoring/service/MonitoringService.java index de5a876..cb99a42 100644 --- a/src/main/java/ru/no_copy/monitoring/service/MonitoringService.java +++ b/src/main/java/ru/no_copy/monitoring/service/MonitoringService.java @@ -1,10 +1,51 @@ package ru.no_copy.monitoring.service; +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Service; +import ru.no_copy.monitoring.dto.MonitoringDTO; +import ru.no_copy.monitoring.dto.SearchResponse; +import ru.no_copy.monitoring.searcher.SearchImageService; + +import java.io.IOException; +import java.util.List; +import java.util.concurrent.TimeoutException; @Service +@RequiredArgsConstructor +@Slf4j public class MonitoringService { - public String execute(String command) { - return "{\"status\": \"ok\", \"data\": \"...\"}"; + + private final SearchImageService imageService; + + private final ObjectMapper objectMapper; + + private final KafkaTemplate kafkaTemplate; + + public void execute(String message) throws IOException, TimeoutException { + MonitoringDTO monitoring; + try { + monitoring = objectMapper.readValue(message, MonitoringDTO.class); + log.info("Monitoring object создан: fileId={}, engine={}", monitoring.getFileId(), monitoring.getEngine()); + } catch (Exception e) { + log.error("Exception readValue: {}", e.getMessage(), e); + throw new RuntimeException(e); + } + + String searchResponse = imageService.searchReverseByPublicUrl(monitoring.getFileId(), monitoring.getBaseUrl(), + monitoring.getEngine(), monitoring.getSearchType()); + log.info("Search response: {}", searchResponse); + + List images = + imageService.getAllImagesWithoutPagination(searchResponse, + monitoring.getSearchType(), + monitoring.getFileId()); + log.info("Images count: {}", images.size()); + + for (SearchResponse.ImageResult imageResult: images) { + kafkaTemplate.send("monitoring-results", imageResult); + } } } diff --git a/src/main/resources/application.yaml b/src/main/resources/application.yaml index ab74f62..01da4c3 100644 --- a/src/main/resources/application.yaml +++ b/src/main/resources/application.yaml @@ -1,13 +1,13 @@ spring: kafka: - bootstrap-servers: kafka:9092 + bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:kafka:9092} consumer: group-id: monitoring-service key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer - value-serializer: org.apache.kafka.common.serialization.StringSerializer + value-serializer: org.springframework.kafka.support.serializer.JsonSerializer cloud: compatibility-verifier: @@ -29,6 +29,10 @@ spring: server: port: ${SERVER_PORT:8083} +searchapi: + api-key: ${SEARCHAPI_API_KEY:5jyYZC8jSaxhZTwjMUhwtAXi} + reverse-image-url: "https://searchapi.io/api/v1/search" + app: internal-api-key: "tljzkXiEYF1klSHuG2hPZKjx6EsBX8RQP6UrzMdQanSKbRYuHuOBwrXejZkn7V4FICvIahoDmYD2hjNBPw61NFbGIt4scOzwZpyCiXEa1YKLAeJSPso4S43LIZlKjO4S"