diff --git a/hub-router/scripts/execution-report.txt b/hub-router/scripts/execution-report.txt new file mode 100644 index 0000000..e808e87 --- /dev/null +++ b/hub-router/scripts/execution-report.txt @@ -0,0 +1,151 @@ +-------------------------------- +Лог выполнения от 22/11/2025 - 20:15:35.200 MSK +-------------------------------- +22/11/2025 - 20:15:35.201 MSK -- Проверяю конфигурацию +22/11/2025 - 20:15:35.203 MSK -- Конфигурация проверена +-------------------------------- +22/11/2025 - 20:15:35.482 MSK -- Проигрываю тестовую сцену: [Проверка обработки событий от устройств и выполнения сценариев умного дома] +-------------------------------- +22/11/2025 - 20:15:35.485 MSK -- Отправляю событие DEVICE_ADDED от хаба hub-1 +22/11/2025 - 20:15:36.505 MSK -- Ожидаю 5 секунд получения сообщения DEVICE_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.084 MSK -- Событие DEVICE_ADDED от хаба hub-1 проверено +-------------------------------- +22/11/2025 - 20:15:38.084 MSK -- Отправляю событие DEVICE_ADDED от хаба hub-1 +22/11/2025 - 20:15:38.090 MSK -- Ожидаю 5 секунд получения сообщения DEVICE_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.100 MSK -- Событие DEVICE_ADDED от хаба hub-1 проверено +-------------------------------- +22/11/2025 - 20:15:38.100 MSK -- Отправляю событие DEVICE_ADDED от хаба hub-1 +22/11/2025 - 20:15:38.106 MSK -- Ожидаю 5 секунд получения сообщения DEVICE_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.110 MSK -- Событие DEVICE_ADDED от хаба hub-1 проверено +-------------------------------- +22/11/2025 - 20:15:38.110 MSK -- Отправляю событие DEVICE_ADDED от хаба hub-1 +22/11/2025 - 20:15:38.113 MSK -- Ожидаю 5 секунд получения сообщения DEVICE_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.118 MSK -- Событие DEVICE_ADDED от хаба hub-1 проверено +-------------------------------- +22/11/2025 - 20:15:38.118 MSK -- Отправляю событие DEVICE_ADDED от хаба hub-1 +22/11/2025 - 20:15:38.121 MSK -- Ожидаю 5 секунд получения сообщения DEVICE_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.126 MSK -- Событие DEVICE_ADDED от хаба hub-1 проверено +-------------------------------- +22/11/2025 - 20:15:38.126 MSK -- Отправляю событие SCENARIO_ADDED от хаба hub-1 +22/11/2025 - 20:15:38.150 MSK -- Ожидаю 5 секунд получения сообщения SCENARIO_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.161 MSK -- Событие SCENARIO_ADDED от хаба hub-1 проверено +-------------------------------- +22/11/2025 - 20:15:38.161 MSK -- Отправляю событие SCENARIO_ADDED от хаба hub-1 +22/11/2025 - 20:15:38.165 MSK -- Ожидаю 5 секунд получения сообщения SCENARIO_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.170 MSK -- Событие SCENARIO_ADDED от хаба hub-1 проверено +-------------------------------- +22/11/2025 - 20:15:38.170 MSK -- Отправляю событие DEVICE_ADDED от хаба hub-2 +22/11/2025 - 20:15:38.174 MSK -- Ожидаю 5 секунд получения сообщения DEVICE_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.179 MSK -- Событие DEVICE_ADDED от хаба hub-2 проверено +-------------------------------- +22/11/2025 - 20:15:38.179 MSK -- Отправляю событие DEVICE_ADDED от хаба hub-2 +22/11/2025 - 20:15:38.182 MSK -- Ожидаю 5 секунд получения сообщения DEVICE_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.187 MSK -- Событие DEVICE_ADDED от хаба hub-2 проверено +-------------------------------- +22/11/2025 - 20:15:38.187 MSK -- Отправляю событие DEVICE_ADDED от хаба hub-2 +22/11/2025 - 20:15:38.191 MSK -- Ожидаю 5 секунд получения сообщения DEVICE_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.196 MSK -- Событие DEVICE_ADDED от хаба hub-2 проверено +-------------------------------- +22/11/2025 - 20:15:38.196 MSK -- Отправляю событие DEVICE_ADDED от хаба hub-2 +22/11/2025 - 20:15:38.199 MSK -- Ожидаю 5 секунд получения сообщения DEVICE_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.207 MSK -- Событие DEVICE_ADDED от хаба hub-2 проверено +-------------------------------- +22/11/2025 - 20:15:38.207 MSK -- Отправляю событие SCENARIO_ADDED от хаба hub-2 +22/11/2025 - 20:15:38.211 MSK -- Ожидаю 5 секунд получения сообщения SCENARIO_ADDED из топика telemetry.hubs.v1 +22/11/2025 - 20:15:38.217 MSK -- Событие SCENARIO_ADDED от хаба hub-2 проверено +-------------------------------- +22/11/2025 - 20:15:38.218 MSK -- Проверяю сценарий [Регулировка температуры (спальня)] для хаба hub-1 +-------------------------------- +22/11/2025 - 20:15:38.218 MSK -- Отправляю событие CLIMATE_SENSOR_EVENT от датчика 2b0bb4c1-7cf2-475a-a17c-e5cb6239d6e5 +22/11/2025 - 20:15:38.252 MSK -- Ожидаю 5 секунд получения сообщения CLIMATE_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.279 MSK -- Получено сообщение CLIMATE_SENSOR_EVENT из топика telemetry.sensors.v1 +-------------------------------- +22/11/2025 - 20:15:38.283 MSK -- Проверяю сценарий [Автосвет (коридор)] для хаба hub-1 +-------------------------------- +22/11/2025 - 20:15:38.283 MSK -- Отправляю событие LIGHT_SENSOR_EVENT от датчика ed9e9587-4148-4fb5-81e0-61d072568628 +22/11/2025 - 20:15:38.292 MSK -- Ожидаю 5 секунд получения сообщения LIGHT_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.298 MSK -- Получено сообщение LIGHT_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.300 MSK -- Отправляю событие MOTION_SENSOR_EVENT от датчика c7b8d4a1-8e37-4c1d-9130-5b0150e13954 +22/11/2025 - 20:15:38.309 MSK -- Ожидаю 5 секунд получения сообщения MOTION_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.314 MSK -- Получено сообщение MOTION_SENSOR_EVENT из топика telemetry.sensors.v1 +-------------------------------- +22/11/2025 - 20:15:38.315 MSK -- Проверяю сценарий [Выключить весь свет] для хаба hub-2 +-------------------------------- +22/11/2025 - 20:15:38.315 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика 14276dbc-d980-4dcc-853e-845ec38c5154 +22/11/2025 - 20:15:38.324 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.330 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +-------------------------------- +22/11/2025 - 20:15:38.333 MSK -- Проверяю пайплайн на случайном наборе событий +-------------------------------- +22/11/2025 - 20:15:38.333 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика 0b9e1641-1a9f-4c43-9b24-6c3f0ccb000e +22/11/2025 - 20:15:38.337 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.343 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.343 MSK -- Отправляю событие LIGHT_SENSOR_EVENT от датчика ed9e9587-4148-4fb5-81e0-61d072568628 +22/11/2025 - 20:15:38.346 MSK -- Ожидаю 5 секунд получения сообщения LIGHT_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.351 MSK -- Получено сообщение LIGHT_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.351 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика 006b61ad-cac8-4adf-9dce-43892b5d060f +22/11/2025 - 20:15:38.354 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.359 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.359 MSK -- Отправляю событие CLIMATE_SENSOR_EVENT от датчика 2b0bb4c1-7cf2-475a-a17c-e5cb6239d6e5 +22/11/2025 - 20:15:38.362 MSK -- Ожидаю 5 секунд получения сообщения CLIMATE_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.365 MSK -- Получено сообщение CLIMATE_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.366 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика b2ec7f40-9c46-4b0b-8ee1-12d0bebde5c1 +22/11/2025 - 20:15:38.368 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.373 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.373 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика 14276dbc-d980-4dcc-853e-845ec38c5154 +22/11/2025 - 20:15:38.377 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.381 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.381 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика 90811ee6-accb-401b-8442-8ec945bdaf29 +22/11/2025 - 20:15:38.384 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.389 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.389 MSK -- Отправляю событие MOTION_SENSOR_EVENT от датчика c7b8d4a1-8e37-4c1d-9130-5b0150e13954 +22/11/2025 - 20:15:38.392 MSK -- Ожидаю 5 секунд получения сообщения MOTION_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.396 MSK -- Получено сообщение MOTION_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.396 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика b2ec7f40-9c46-4b0b-8ee1-12d0bebde5c1 +22/11/2025 - 20:15:38.399 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.405 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.405 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика b2ec7f40-9c46-4b0b-8ee1-12d0bebde5c1 +22/11/2025 - 20:15:38.408 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.413 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.414 MSK -- Отправляю событие CLIMATE_SENSOR_EVENT от датчика 2b0bb4c1-7cf2-475a-a17c-e5cb6239d6e5 +22/11/2025 - 20:15:38.416 MSK -- Ожидаю 5 секунд получения сообщения CLIMATE_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.422 MSK -- Получено сообщение CLIMATE_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.422 MSK -- Отправляю событие LIGHT_SENSOR_EVENT от датчика ed9e9587-4148-4fb5-81e0-61d072568628 +22/11/2025 - 20:15:38.425 MSK -- Ожидаю 5 секунд получения сообщения LIGHT_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.430 MSK -- Получено сообщение LIGHT_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.430 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика 006b61ad-cac8-4adf-9dce-43892b5d060f +22/11/2025 - 20:15:38.433 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.438 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.438 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика 0b9e1641-1a9f-4c43-9b24-6c3f0ccb000e +22/11/2025 - 20:15:38.442 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.447 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.447 MSK -- Отправляю событие MOTION_SENSOR_EVENT от датчика c7b8d4a1-8e37-4c1d-9130-5b0150e13954 +22/11/2025 - 20:15:38.451 MSK -- Ожидаю 5 секунд получения сообщения MOTION_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.457 MSK -- Получено сообщение MOTION_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.457 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика 006b61ad-cac8-4adf-9dce-43892b5d060f +22/11/2025 - 20:15:38.460 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.465 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.465 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика 0b9e1641-1a9f-4c43-9b24-6c3f0ccb000e +22/11/2025 - 20:15:38.468 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.474 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.474 MSK -- Отправляю событие LIGHT_SENSOR_EVENT от датчика ed9e9587-4148-4fb5-81e0-61d072568628 +22/11/2025 - 20:15:38.477 MSK -- Ожидаю 5 секунд получения сообщения LIGHT_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.482 MSK -- Получено сообщение LIGHT_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.482 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика f94d84a9-c9dd-41df-bbdc-70a8e609437b +22/11/2025 - 20:15:38.485 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.490 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.490 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика 14276dbc-d980-4dcc-853e-845ec38c5154 +22/11/2025 - 20:15:38.493 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.497 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.498 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика 90811ee6-accb-401b-8442-8ec945bdaf29 +22/11/2025 - 20:15:38.501 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.505 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.505 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика b2ec7f40-9c46-4b0b-8ee1-12d0bebde5c1 +22/11/2025 - 20:15:38.509 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.512 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.513 MSK -- Отправляю событие SWITCH_SENSOR_EVENT от датчика f94d84a9-c9dd-41df-bbdc-70a8e609437b +22/11/2025 - 20:15:38.516 MSK -- Ожидаю 5 секунд получения сообщения SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +22/11/2025 - 20:15:38.520 MSK -- Получено сообщение SWITCH_SENSOR_EVENT из топика telemetry.sensors.v1 +-------------------------------- +-------------------------------- +Обнаружено 0 ошибок. diff --git a/hub-router/scripts/windows/2-collector-grpc-tests.bat b/hub-router/scripts/windows/2-collector-grpc-tests.bat new file mode 100644 index 0000000..d82535f --- /dev/null +++ b/hub-router/scripts/windows/2-collector-grpc-tests.bat @@ -0,0 +1,30 @@ +@echo off +setlocal + +REM Локальный запуск тестов Hub Router для проверки сервиса Collector + +set "JAR_PATH=%~dp0..\hub-router.jar" + +if "%1"=="info" ( + echo. + java -jar "%JAR_PATH%" info + echo. + pause + exit /b +) + +echo "Запуск Hub Router (режим: COLLECTION, GRPC)" +echo. + +java -jar "%JAR_PATH%" ^ + --hub-router.execution.mode=COLLECTION ^ + --hub-router.execution.collector.mode=grpc ^ + --hub-router.execution.immediate-logging.enabled=false ^ + --hub-router.execution.output.info-enabled=true ^ + --hub-router.execution.output.trace-enabled=true ^ + --hub-router.execution.output.console=true ^ + --hub-router.skip-summary-on-startup=false + +echo. +echo Тест завершён. Проверьте результаты в консоли выше. +pause \ No newline at end of file diff --git a/hub-router/start-hub-router.bat b/hub-router/start-hub-router.bat new file mode 100644 index 0000000..1efdf32 --- /dev/null +++ b/hub-router/start-hub-router.bat @@ -0,0 +1,36 @@ +@echo off +setlocal + +if "%~1"=="" ( + echo ❌ Ошибка: Не передан аргумент! Укажите один из: 1-collector-json, 2-collector-grpc, 3-aggregator, 4-analyzer + exit /b 1 +) + +set MODE=%~1 + +if "%MODE%"=="1-collector-json" ( + echo 🚀 Запуск hub-router в режиме HTTP Collector... + java -jar hub-router.jar --hub-router.execution.collector.mode=http --hub-router.execution.collector.port=8080 + exit /b +) + +if "%MODE%"=="2-collector-grpc" ( + echo 🚀 Запуск hub-router в режиме gRPC Collector... + java -jar hub-router.jar + exit /b +) + +if "%MODE%"=="3-aggregator" ( + echo 🚀 Запуск hub-router в режиме Aggregator... + java -jar hub-router.jar --hub-router.execution.mode=AGGREGATION + exit /b +) + +if "%MODE%"=="4-analyzer" ( + echo 🚀 Запуск hub-router в режиме Analyzer... + java -jar hub-router.jar --hub-router.execution.mode=ANALYZE + exit /b +) + +echo ❌ Ошибка: Неверный аргумент '%MODE%'. Доступные варианты: 1-collector-json, 2-collector-grpc, 3-aggregator, 4-analyzer +exit /b 1 \ No newline at end of file diff --git a/telemetry/collector/pom.xml b/telemetry/collector/pom.xml index 6f54a15..0741d2b 100644 --- a/telemetry/collector/pom.xml +++ b/telemetry/collector/pom.xml @@ -31,12 +31,31 @@ org.projectlombok lombok + true org.springframework.kafka spring-kafka + + + ru.yandex.practicum + proto-schemas + 1.0-SNAPSHOT + + + + net.devh + grpc-server-spring-boot-starter + 3.1.0.RELEASE + + + + com.google.protobuf + protobuf-java-util + ${protobuf.version} + @@ -44,13 +63,29 @@ org.springframework.boot spring-boot-maven-plugin + + + + repackage + + + + + + + org.apache.maven.plugins + maven-compiler-plugin + 3.11.0 - - + 21 + 21 + + org.projectlombok lombok - - + ${lombok.version} + + diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/controller/CollectorGrpcController.java b/telemetry/collector/src/main/java/ru/yandex/practicum/controller/CollectorGrpcController.java new file mode 100644 index 0000000..688d4b3 --- /dev/null +++ b/telemetry/collector/src/main/java/ru/yandex/practicum/controller/CollectorGrpcController.java @@ -0,0 +1,57 @@ +package ru.yandex.practicum.controller; + +import com.google.protobuf.Empty; +import io.grpc.Status; +import io.grpc.StatusRuntimeException; +import io.grpc.stub.StreamObserver; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import net.devh.boot.grpc.server.service.GrpcService; +import ru.yandex.practicum.grpc.telemetry.collector.CollectorControllerGrpc; +import ru.yandex.practicum.grpc.telemetry.event.HubEventProto; +import ru.yandex.practicum.grpc.telemetry.event.SensorEventProto; +import ru.yandex.practicum.mapper.*; +import ru.yandex.practicum.service.EventService; + +@Slf4j +@GrpcService +@RequiredArgsConstructor +public class CollectorGrpcController extends CollectorControllerGrpc.CollectorControllerImplBase { + + private final EventService eventService; + private final ProtoToAvroSensorMapper sensorMapper; + private final ProtoToAvroHubMapper hubMapper; + + @Override + public void collectSensorEvent(SensorEventProto request, StreamObserver responseObserver) { + log.info("gRPC: получен SensorEventProto: {}", request); + try { + var avro = sensorMapper.toAvro(request); + eventService.sendSensorEvent(avro); + responseObserver.onNext(Empty.getDefaultInstance()); + responseObserver.onCompleted(); + } catch (Exception e) { + handleError(responseObserver, e, "collectSensorEvent"); + } + } + + @Override + public void collectHubEvent(HubEventProto request, StreamObserver responseObserver) { + log.info("gRPC: получен HubEventProto: {}", request); + try { + var avro = hubMapper.toAvro(request); + eventService.sendHubEvent(avro); + responseObserver.onNext(Empty.getDefaultInstance()); + responseObserver.onCompleted(); + } catch (Exception e) { + handleError(responseObserver, e, "collectHubEvent"); + } + } + + private void handleError(StreamObserver responseObserver, Exception e, String context) { + log.error("Ошибка в {}: {}", context, e.getMessage(), e); + responseObserver.onError(new StatusRuntimeException( + Status.INTERNAL.withDescription(e.getLocalizedMessage()).withCause(e) + )); + } +} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/controller/HubEventController.java b/telemetry/collector/src/main/java/ru/yandex/practicum/controller/HubEventController.java deleted file mode 100644 index 1a51463..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/controller/HubEventController.java +++ /dev/null @@ -1,39 +0,0 @@ -package ru.yandex.practicum.controller; - -import jakarta.validation.Valid; -import lombok.RequiredArgsConstructor; -import lombok.extern.slf4j.Slf4j; -import org.springframework.http.ResponseEntity; -import org.springframework.web.bind.annotation.PostMapping; -import org.springframework.web.bind.annotation.RequestBody; -import org.springframework.web.bind.annotation.RequestMapping; -import org.springframework.web.bind.annotation.RestController; -import ru.yandex.practicum.mapper.HubEventMapper; -import ru.yandex.practicum.mapper.SensorEventMapper; -import ru.yandex.practicum.model.hub.HubEvent; -import ru.yandex.practicum.model.sensor.SensorEvent; -import ru.yandex.practicum.service.EventService; - -@Slf4j -@RestController -@RequestMapping("/events") -@RequiredArgsConstructor -public class HubEventController { - private final SensorEventMapper sensorEventMapper; - private final HubEventMapper hubEventMapper; - private final EventService service; - - @PostMapping("/sensors") - public ResponseEntity collectSensor(@Valid @RequestBody SensorEvent sensorEvent) { - log.info("Received SensorEvent: {}", sensorEvent); - service.sendSensorEvent(sensorEventMapper.toAvro(sensorEvent)); - return ResponseEntity.ok().build(); - } - - @PostMapping("/hubs") - public ResponseEntity collectHub(@Valid @RequestBody HubEvent hubEvent) { - log.info("Received HubEvent: {}", hubEvent); - service.sendHubEvent(hubEventMapper.toAvro(hubEvent)); - return ResponseEntity.ok().build(); - } -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/HubEventMapper.java b/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/HubEventMapper.java deleted file mode 100644 index d47badc..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/HubEventMapper.java +++ /dev/null @@ -1,77 +0,0 @@ -package ru.yandex.practicum.mapper; - -import lombok.extern.slf4j.Slf4j; -import org.springframework.stereotype.Component; -import ru.yandex.practicum.kafka.telemetry.event.*; -import ru.yandex.practicum.model.hub.*; - -import java.util.stream.Collectors; - -@Slf4j -@Component -public class HubEventMapper { - public HubEventAvro toAvro(HubEvent event) { - log.info("Mapping HubEvent type: {}", event.getType()); - HubEventAvro.Builder builder = HubEventAvro.newBuilder() - .setHubId(event.getHubId()) - .setTimestamp(event.getTimestamp()); - - switch (event.getType()) { - case DEVICE_ADDED -> { - DeviceAddedEvent added = (DeviceAddedEvent) event; - log.debug("Mapping DEVICE_ADDED event - id: {}, deviceType: {}", added.getId(), added.getDeviceType()); - builder.setPayload( - DeviceAddedEventAvro.newBuilder() - .setId(added.getId()) - .setType(DeviceTypeAvro.valueOf(added.getDeviceType().name())) - .build() - ); - } - - case DEVICE_REMOVED -> { - DeviceRemovedEvent removed = (DeviceRemovedEvent) event; - log.debug("Mapping DEVICE_REMOVED event - id: {}", removed.getId()); - builder.setPayload( - DeviceRemovedEventAvro.newBuilder() - .setId(removed.getId()) - .build() - ); - } - case SCENARIO_ADDED -> { - ScenarioAddedEvent added = (ScenarioAddedEvent) event; - log.debug("Mapping SCENARIO_ADDED event - name: {}, conditions: {}, actions: {}", - added.getName(), added.getConditions().size(), added.getActions().size()); - builder.setPayload( - ScenarioAddedEventAvro.newBuilder() - .setName(added.getName()) - .setConditions(added.getConditions().stream() - .map(c -> ScenarioConditionAvro.newBuilder() - .setSensorId(c.getSensorId()) - .setType(ConditionTypeAvro.valueOf(c.getType().name())) - .setOperation(ConditionOperationAvro.valueOf(c.getOperation().name())) - .setValue(c.getValue()) - .build()) - .collect(Collectors.toList())) - .setActions(added.getActions().stream() - .map(a -> DeviceActionAvro.newBuilder() - .setSensorId(a.getSensorId()) - .setType(ActionTypeAvro.valueOf(a.getType().name())) - .setValue(a.getValue()) - .build()) - .collect(Collectors.toList())) - .build() - ); - } - case SCENARIO_REMOVED -> { - ScenarioRemovedEvent removed = (ScenarioRemovedEvent) event; - log.debug("Mapping SCENARIO_REMOVED event - name: {}", removed.getName()); - builder.setPayload( - ScenarioRemovedEventAvro.newBuilder() - .setName(removed.getName()) - .build() - ); - } - } - return builder.build(); - } -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/ProtoToAvroHubMapper.java b/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/ProtoToAvroHubMapper.java new file mode 100644 index 0000000..b6526ea --- /dev/null +++ b/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/ProtoToAvroHubMapper.java @@ -0,0 +1,93 @@ +package ru.yandex.practicum.mapper; + +import org.springframework.stereotype.Component; +import ru.yandex.practicum.grpc.telemetry.event.*; +import ru.yandex.practicum.kafka.telemetry.event.*; + +import java.time.Instant; +import java.util.stream.Collectors; + +@Component +public class ProtoToAvroHubMapper { + public HubEventAvro toAvro(HubEventProto proto) { + HubEventAvro.Builder builder = HubEventAvro.newBuilder() + .setHubId(proto.getHubId()) + .setTimestamp(Instant.ofEpochSecond( + proto.getTimestamp().getSeconds(), + proto.getTimestamp().getNanos())); + + return switch (proto.getPayloadCase()) { + case DEVICE_ADDED -> { + DeviceAddedEventProto d = proto.getDeviceAdded(); + yield builder.setPayload( + DeviceAddedEventAvro.newBuilder() + .setId(d.getId()) + .setType(DeviceTypeAvro.valueOf(d.getType().name())) + .build() + ).build(); + } + case DEVICE_REMOVED -> { + DeviceRemovedEventProto d = proto.getDeviceRemoved(); + yield builder.setPayload( + DeviceRemovedEventAvro.newBuilder() + .setId(d.getId()) + .build() + ).build(); + } + case SCENARIO_ADDED -> { + ScenarioAddedEventProto s = proto.getScenarioAdded(); + yield builder.setPayload( + ScenarioAddedEventAvro.newBuilder() + .setName(s.getName()) + .setConditions( + s.getConditionList().stream() + .map(this::mapCondition) + .collect(Collectors.toList()) + ) + .setActions( + s.getActionList().stream() + .map(this::mapAction) + .collect(Collectors.toList()) + ) + .build() + ).build(); + } + case SCENARIO_REMOVED -> { + ScenarioRemovedEventProto s = proto.getScenarioRemoved(); + yield builder.setPayload( + ScenarioRemovedEventAvro.newBuilder() + .setName(s.getName()) + .build() + ).build(); + } + case PAYLOAD_NOT_SET -> builder.build(); + }; + } + + private ScenarioConditionAvro mapCondition(ScenarioConditionProto proto) { + ScenarioConditionAvro.Builder b = ScenarioConditionAvro.newBuilder() + .setSensorId(proto.getSensorId()) + .setType(ConditionTypeAvro.valueOf(proto.getType().name())) + .setOperation(ConditionOperationAvro.valueOf(proto.getOperation().name())); + + switch (proto.getValueCase()) { + case BOOL_VALUE -> b.setValue(proto.getBoolValue()); + case INT_VALUE -> b.setValue(proto.getIntValue()); + case VALUE_NOT_SET -> b.setValue(null); + } + return b.build(); + } + + private DeviceActionAvro mapAction(DeviceActionProto proto) { + DeviceActionAvro.Builder b = DeviceActionAvro.newBuilder() + .setSensorId(proto.getSensorId()) + .setType(ActionTypeAvro.valueOf(proto.getType().name())); + + if (proto.hasValue()) { + b.setValue(proto.getValue()); + } else { + b.setValue(null); + } + return b.build(); + } +} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/ProtoToAvroSensorMapper.java b/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/ProtoToAvroSensorMapper.java new file mode 100644 index 0000000..0a6cf48 --- /dev/null +++ b/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/ProtoToAvroSensorMapper.java @@ -0,0 +1,65 @@ +package ru.yandex.practicum.mapper; + +import org.springframework.stereotype.Component; +import ru.yandex.practicum.grpc.telemetry.event.*; +import ru.yandex.practicum.kafka.telemetry.event.*; + +import java.time.Instant; + +@Component +public class ProtoToAvroSensorMapper { + + public SensorEventAvro toAvro(SensorEventProto proto) { + SensorEventAvro.Builder builder = SensorEventAvro.newBuilder() + .setId(proto.getId()) + .setHubId(proto.getHubId()) + .setTimestamp(Instant.ofEpochSecond( + proto.getTimestamp().getSeconds(), + proto.getTimestamp().getNanos() + )); + return switch (proto.getPayloadCase()) { + case MOTION_SENSOR -> { + MotionSensorProto m = proto.getMotionSensor(); + yield builder.setPayload(MotionSensorAvro.newBuilder() + .setLinkQuality(m.getLinkQuality()) + .setMotion(m.getMotion()) + .setVoltage(m.getVoltage()) + .build()) + .build(); + } + case TEMPERATURE_SENSOR -> { + TemperatureSensorProto t = proto.getTemperatureSensor(); + yield builder.setPayload(TemperatureSensorAvro.newBuilder() + .setTemperatureC(t.getTemperatureC()) + .setTemperatureF(t.getTemperatureF()) + .build()) + .build(); + } + case LIGHT_SENSOR -> { + LightSensorProto l = proto.getLightSensor(); + yield builder.setPayload(LightSensorAvro.newBuilder() + .setLinkQuality(l.getLinkQuality()) + .setLuminosity(l.getLuminosity()) + .build()) + .build(); + } + case CLIMATE_SENSOR -> { + ClimateSensorProto c = proto.getClimateSensor(); + yield builder.setPayload(ClimateSensorAvro.newBuilder() + .setTemperatureC(c.getTemperatureC()) + .setHumidity(c.getHumidity()) + .setCo2Level(c.getCo2Level()) + .build()) + .build(); + } + case SWITCH_SENSOR -> { + SwitchSensorProto s = proto.getSwitchSensor(); + yield builder.setPayload(SwitchSensorAvro.newBuilder() + .setState(s.getState()) + .build()) + .build(); + } + case PAYLOAD_NOT_SET -> builder.build(); + }; + } +} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/SensorEventMapper.java b/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/SensorEventMapper.java deleted file mode 100644 index 6b445d5..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/mapper/SensorEventMapper.java +++ /dev/null @@ -1,50 +0,0 @@ -package ru.yandex.practicum.mapper; - -import org.springframework.stereotype.Component; -import ru.yandex.practicum.kafka.telemetry.event.*; -import ru.yandex.practicum.model.sensor.*; - -@Component -public class SensorEventMapper { - public SensorEventAvro toAvro(SensorEvent event) { - SensorEventAvro.Builder builder = SensorEventAvro.newBuilder() - .setId(event.getId()) - .setHubId(event.getHubId()) - .setTimestamp(event.getTimestamp()); - - switch (event.getType()) { - case LIGHT_SENSOR_EVENT -> builder.setPayload( - LightSensorAvro.newBuilder() - .setLinkQuality(((LightSensorEvent) event).getLinkQuality()) - .setLuminosity(((LightSensorEvent) event).getLuminosity()) - .build()); - - case MOTION_SENSOR_EVENT -> builder.setPayload( - MotionSensorAvro.newBuilder() - .setLinkQuality(((MotionSensorEvent) event).getLinkQuality()) - .setMotion(((MotionSensorEvent) event).isMotion()) - .setVoltage((((MotionSensorEvent) event).getVoltage())) - .build()); - - case CLIMATE_SENSOR_EVENT -> builder.setPayload( - ClimateSensorAvro.newBuilder() - .setTemperatureC(((ClimateSensorEvent) event).getTemperatureC()) - .setHumidity(((ClimateSensorEvent) event).getHumidity()) - .setCo2Level(((ClimateSensorEvent) event).getCo2Level()) - .build()); - - case SWITCH_SENSOR_EVENT -> builder.setPayload( - SwitchSensorAvro.newBuilder() - .setState(((SwitchSensorEvent) event).isState()) - .build()); - - case TEMPERATURE_SENSOR_EVENT -> builder.setPayload( - TemperatureSensorAvro.newBuilder() - .setTemperatureC(((TemperatureSensorEvent) event).getTemperatureC()) - .setTemperatureF(((TemperatureSensorEvent) event).getTemperatureF()) - .build()); - } - - return builder.build(); - } -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ActionType.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ActionType.java deleted file mode 100644 index 64fa33d..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ActionType.java +++ /dev/null @@ -1,8 +0,0 @@ -package ru.yandex.practicum.model.hub; - -public enum ActionType { - ACTIVATE, - DEACTIVATE, - INVERSE, - SET_VALUE -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ConditionOperation.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ConditionOperation.java deleted file mode 100644 index 741c433..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ConditionOperation.java +++ /dev/null @@ -1,7 +0,0 @@ -package ru.yandex.practicum.model.hub; - -public enum ConditionOperation { - EQUALS, - GREATER_THAN, - LOWER_THAN -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ConditionType.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ConditionType.java deleted file mode 100644 index 57ba1fc..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ConditionType.java +++ /dev/null @@ -1,10 +0,0 @@ -package ru.yandex.practicum.model.hub; - -public enum ConditionType { - MOTION, - LUMINOSITY, - SWITCH, - TEMPERATURE, - CO2LEVEL, - HUMIDITY -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceAction.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceAction.java deleted file mode 100644 index 4b3e292..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceAction.java +++ /dev/null @@ -1,19 +0,0 @@ -package ru.yandex.practicum.model.hub; - -import jakarta.validation.constraints.NotBlank; -import jakarta.validation.constraints.NotNull; -import lombok.Getter; -import lombok.Setter; - -@Getter -@Setter -public class DeviceAction { - - @NotBlank - private String sensorId; - - @NotNull - private ActionType type; - - private Integer value; -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceAddedEvent.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceAddedEvent.java deleted file mode 100644 index 01a9e2c..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceAddedEvent.java +++ /dev/null @@ -1,22 +0,0 @@ -package ru.yandex.practicum.model.hub; - -import jakarta.validation.constraints.NotBlank; -import jakarta.validation.constraints.NotNull; -import lombok.Getter; -import lombok.Setter; - -@Getter -@Setter -public class DeviceAddedEvent extends HubEvent { - - @NotBlank - private String id; - - @NotNull - private DeviceType deviceType; - - @Override - public HubEventType getType() { - return HubEventType.DEVICE_ADDED; - } -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceRemovedEvent.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceRemovedEvent.java deleted file mode 100644 index ec463ae..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceRemovedEvent.java +++ /dev/null @@ -1,18 +0,0 @@ -package ru.yandex.practicum.model.hub; - -import jakarta.validation.constraints.NotBlank; -import lombok.Getter; -import lombok.Setter; - -@Getter -@Setter -public class DeviceRemovedEvent extends HubEvent { - - @NotBlank - String id; - - @Override - public HubEventType getType() { - return HubEventType.DEVICE_REMOVED; - } -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceType.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceType.java deleted file mode 100644 index e01d764..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/DeviceType.java +++ /dev/null @@ -1,9 +0,0 @@ -package ru.yandex.practicum.model.hub; - -public enum DeviceType { - MOTION_SENSOR, - TEMPERATURE_SENSOR, - LIGHT_SENSOR, - CLIMATE_SENSOR, - SWITCH_SENSOR -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/HubEvent.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/HubEvent.java deleted file mode 100644 index 1f614ea..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/HubEvent.java +++ /dev/null @@ -1,37 +0,0 @@ -package ru.yandex.practicum.model.hub; - -import com.fasterxml.jackson.annotation.JsonSubTypes; -import com.fasterxml.jackson.annotation.JsonTypeInfo; -import jakarta.validation.constraints.NotBlank; -import jakarta.validation.constraints.NotNull; -import lombok.Getter; -import lombok.Setter; -import lombok.ToString; - -import java.time.Instant; - -@JsonTypeInfo( - use = JsonTypeInfo.Id.NAME, - include = JsonTypeInfo.As.EXISTING_PROPERTY, - property = "type" -) -@JsonSubTypes({ - @JsonSubTypes.Type(value = DeviceAddedEvent.class, name = "DEVICE_ADDED"), - @JsonSubTypes.Type(value = DeviceRemovedEvent.class, name = "DEVICE_REMOVED"), - @JsonSubTypes.Type(value = ScenarioAddedEvent.class, name = "SCENARIO_ADDED"), - @JsonSubTypes.Type(value = ScenarioRemovedEvent.class, name = "SCENARIO_REMOVED") -}) - -@Getter -@Setter -@ToString -public abstract class HubEvent { - - @NotBlank - private String hubId; - - private Instant timestamp = Instant.now(); - - @NotNull - public abstract HubEventType getType(); -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/HubEventType.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/HubEventType.java deleted file mode 100644 index a3188d7..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/HubEventType.java +++ /dev/null @@ -1,8 +0,0 @@ -package ru.yandex.practicum.model.hub; - -public enum HubEventType { - DEVICE_ADDED, - DEVICE_REMOVED, - SCENARIO_ADDED, - SCENARIO_REMOVED -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ScenarioAddedEvent.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ScenarioAddedEvent.java deleted file mode 100644 index 1ad0d8a..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ScenarioAddedEvent.java +++ /dev/null @@ -1,27 +0,0 @@ -package ru.yandex.practicum.model.hub; - -import jakarta.validation.constraints.NotBlank; -import jakarta.validation.constraints.NotNull; -import lombok.Getter; -import lombok.Setter; - -import java.util.List; - -@Getter -@Setter -public class ScenarioAddedEvent extends HubEvent { - - @NotBlank - private String name; - - @NotNull - private List conditions; - - @NotNull - private List actions; - - @Override - public HubEventType getType() { - return HubEventType.SCENARIO_ADDED; - } -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ScenarioCondition.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ScenarioCondition.java deleted file mode 100644 index 2ead489..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ScenarioCondition.java +++ /dev/null @@ -1,23 +0,0 @@ -package ru.yandex.practicum.model.hub; - -import jakarta.validation.constraints.NotBlank; -import jakarta.validation.constraints.NotNull; -import lombok.Getter; -import lombok.Setter; - -@Getter -@Setter -public class ScenarioCondition { - - @NotBlank - private String sensorId; - - @NotNull - private ConditionType type; - - @NotNull - private ConditionOperation operation; - - @NotNull - private Integer value; -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ScenarioRemovedEvent.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ScenarioRemovedEvent.java deleted file mode 100644 index 45b55b7..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/hub/ScenarioRemovedEvent.java +++ /dev/null @@ -1,18 +0,0 @@ -package ru.yandex.practicum.model.hub; - -import jakarta.validation.constraints.NotBlank; -import lombok.Getter; -import lombok.Setter; - -@Getter -@Setter -public class ScenarioRemovedEvent extends HubEvent { - - @NotBlank - private String name; - - @Override - public HubEventType getType() { - return HubEventType.SCENARIO_REMOVED; - } -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/ClimateSensorEvent.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/ClimateSensorEvent.java deleted file mode 100644 index d70ee49..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/ClimateSensorEvent.java +++ /dev/null @@ -1,24 +0,0 @@ -package ru.yandex.practicum.model.sensor; - -import jakarta.validation.constraints.NotNull; -import lombok.Getter; -import lombok.Setter; - -@Getter -@Setter -public class ClimateSensorEvent extends SensorEvent { - - @NotNull - private Integer temperatureC; - - @NotNull - private Integer humidity; - - @NotNull - private Integer co2Level; - - @Override - public SensorEventType getType() { - return SensorEventType.CLIMATE_SENSOR_EVENT; - } -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/LightSensorEvent.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/LightSensorEvent.java deleted file mode 100644 index 87a985d..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/LightSensorEvent.java +++ /dev/null @@ -1,21 +0,0 @@ -package ru.yandex.practicum.model.sensor; - -import jakarta.validation.constraints.NotNull; -import lombok.Getter; -import lombok.Setter; - -@Getter -@Setter -public class LightSensorEvent extends SensorEvent { - - @NotNull - private Integer linkQuality; - - @NotNull - private Integer luminosity; - - @Override - public SensorEventType getType() { - return SensorEventType.LIGHT_SENSOR_EVENT; - } -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/MotionSensorEvent.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/MotionSensorEvent.java deleted file mode 100644 index 3c1edf3..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/MotionSensorEvent.java +++ /dev/null @@ -1,24 +0,0 @@ -package ru.yandex.practicum.model.sensor; - -import jakarta.validation.constraints.NotNull; -import lombok.Getter; -import lombok.Setter; - -@Getter -@Setter -public class MotionSensorEvent extends SensorEvent { - - @NotNull - private Integer linkQuality; - - @NotNull - private boolean motion; - - @NotNull - private Integer voltage; - - @Override - public SensorEventType getType() { - return SensorEventType.MOTION_SENSOR_EVENT; - } -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/SensorEvent.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/SensorEvent.java deleted file mode 100644 index f151c63..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/SensorEvent.java +++ /dev/null @@ -1,42 +0,0 @@ -package ru.yandex.practicum.model.sensor; - -import com.fasterxml.jackson.annotation.JsonSubTypes; -import com.fasterxml.jackson.annotation.JsonTypeInfo; -import jakarta.validation.constraints.NotBlank; -import jakarta.validation.constraints.NotNull; -import lombok.Getter; -import lombok.Setter; -import lombok.ToString; - -import java.time.Instant; - -@JsonTypeInfo( - use = JsonTypeInfo.Id.NAME, - include = JsonTypeInfo.As.EXISTING_PROPERTY, - property = "type" -) -@JsonSubTypes({ - @JsonSubTypes.Type(value = LightSensorEvent.class, name = "LIGHT_SENSOR_EVENT"), - @JsonSubTypes.Type(value = MotionSensorEvent.class, name = "MOTION_SENSOR_EVENT"), - @JsonSubTypes.Type(value = TemperatureSensorEvent.class, name = "TEMPERATURE_SENSOR_EVENT"), - @JsonSubTypes.Type(value = ClimateSensorEvent.class, name = "CLIMATE_SENSOR_EVENT"), - @JsonSubTypes.Type(value = SwitchSensorEvent.class, name = "SWITCH_SENSOR_EVENT") -}) - -@Getter -@Setter -@ToString -public abstract class SensorEvent { - - @NotBlank - private String id; - - @NotBlank - private String hubId; - - @NotNull - private Instant timestamp; - - @NotNull - public abstract SensorEventType getType(); -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/SensorEventType.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/SensorEventType.java deleted file mode 100644 index e158a93..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/SensorEventType.java +++ /dev/null @@ -1,9 +0,0 @@ -package ru.yandex.practicum.model.sensor; - -public enum SensorEventType { - MOTION_SENSOR_EVENT, - TEMPERATURE_SENSOR_EVENT, - LIGHT_SENSOR_EVENT, - CLIMATE_SENSOR_EVENT, - SWITCH_SENSOR_EVENT -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/SwitchSensorEvent.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/SwitchSensorEvent.java deleted file mode 100644 index 3eff96a..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/SwitchSensorEvent.java +++ /dev/null @@ -1,18 +0,0 @@ -package ru.yandex.practicum.model.sensor; - -import jakarta.validation.constraints.NotNull; -import lombok.Getter; -import lombok.Setter; - -@Getter -@Setter -public class SwitchSensorEvent extends SensorEvent { - - @NotNull - private boolean state; - - @Override - public SensorEventType getType() { - return SensorEventType.SWITCH_SENSOR_EVENT; - } -} diff --git a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/TemperatureSensorEvent.java b/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/TemperatureSensorEvent.java deleted file mode 100644 index 4d27769..0000000 --- a/telemetry/collector/src/main/java/ru/yandex/practicum/model/sensor/TemperatureSensorEvent.java +++ /dev/null @@ -1,21 +0,0 @@ -package ru.yandex.practicum.model.sensor; - -import jakarta.validation.constraints.NotNull; -import lombok.Getter; -import lombok.Setter; - -@Getter -@Setter -public class TemperatureSensorEvent extends SensorEvent { - - @NotNull - private Integer temperatureC; - - @NotNull - private Integer temperatureF; - - @Override - public SensorEventType getType() { - return SensorEventType.TEMPERATURE_SENSOR_EVENT; - } -} diff --git a/telemetry/collector/src/main/resources/application.yaml b/telemetry/collector/src/main/resources/application.yaml index 941cc23..8e3b9ca 100644 --- a/telemetry/collector/src/main/resources/application.yaml +++ b/telemetry/collector/src/main/resources/application.yaml @@ -13,6 +13,11 @@ collector: sensors: telemetry.sensors.v1 hubs: telemetry.hubs.v1 +grpc: + server: + port: 59091 + reflection-service-enabled: true + logging: level: root: INFO diff --git a/telemetry/serialization/proto-schemas/pom.xml b/telemetry/serialization/proto-schemas/pom.xml index 769e58e..f52ac1e 100644 --- a/telemetry/serialization/proto-schemas/pom.xml +++ b/telemetry/serialization/proto-schemas/pom.xml @@ -45,7 +45,6 @@ ${protobuf.version} - io.grpc diff --git a/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/messages/hub_event.proto b/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/messages/hub_event.proto index f978af3..50bf5a9 100644 --- a/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/messages/hub_event.proto +++ b/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/messages/hub_event.proto @@ -3,4 +3,82 @@ syntax = "proto3"; package telemetry.message.event; option java_multiple_files = true; -option java_package = "ru.yandex.practicum.grpc.telemetry.event"; \ No newline at end of file +option java_package = "ru.yandex.practicum.grpc.telemetry.event"; + +import "google/protobuf/timestamp.proto"; + +message HubEventProto { + string hubId = 1; + google.protobuf.Timestamp timestamp = 2; + oneof payload { + DeviceAddedEventProto device_added = 3; + DeviceRemovedEventProto device_removed = 4; + ScenarioAddedEventProto scenario_added = 5; + ScenarioRemovedEventProto scenario_removed = 6; + } +} + +message DeviceAddedEventProto { + string id = 1; + DeviceTypeProto type = 2; +} + +message DeviceRemovedEventProto { + string id = 1; +} + +message ScenarioConditionProto { + string sensor_id = 1; + ConditionTypeProto type = 2; + ConditionOperationProto operation = 3; + oneof value { + bool bool_value = 4; + int32 int_value = 5; + } +} + +message DeviceActionProto { + string sensor_id = 1; + ActionTypeProto type = 2; + optional int32 value = 3; +} + +message ScenarioAddedEventProto { + string name = 1; + repeated ScenarioConditionProto condition = 2; + repeated DeviceActionProto action = 3; +} + +message ScenarioRemovedEventProto { + string name = 1; +} + +enum DeviceTypeProto { + MOTION_SENSOR = 0; + TEMPERATURE_SENSOR = 1; + LIGHT_SENSOR = 2; + CLIMATE_SENSOR = 3; + SWITCH_SENSOR = 4; +} + +enum ConditionTypeProto { + MOTION = 0; + LUMINOSITY = 1; + SWITCH = 2; + TEMPERATURE = 3; + CO2LEVEL = 4; + HUMIDITY = 5; +} + +enum ConditionOperationProto { + EQUALS = 0; + GREATER_THAN = 1; + LOWER_THAN = 2; +} + +enum ActionTypeProto { + ACTIVATE = 0; + DEACTIVATE = 1; + INVERSE = 2; + SET_VALUE = 3; +} diff --git a/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/messages/sensor_event.proto b/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/messages/sensor_event.proto index f978af3..16b3632 100644 --- a/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/messages/sensor_event.proto +++ b/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/messages/sensor_event.proto @@ -3,4 +3,45 @@ syntax = "proto3"; package telemetry.message.event; option java_multiple_files = true; -option java_package = "ru.yandex.practicum.grpc.telemetry.event"; \ No newline at end of file +option java_package = "ru.yandex.practicum.grpc.telemetry.event"; + +import "google/protobuf/timestamp.proto"; + +message SensorEventProto { + string id = 1; + google.protobuf.Timestamp timestamp = 2; + string hubId = 3; + oneof payload { + MotionSensorProto motion_sensor = 4; + TemperatureSensorProto temperature_sensor = 5; + LightSensorProto light_sensor = 6; + ClimateSensorProto climate_sensor = 7; + SwitchSensorProto switch_sensor = 8; + } +} + +message MotionSensorProto { + int32 link_quality = 1; + bool motion = 2; + int32 voltage = 3; +} + +message TemperatureSensorProto { + int32 temperature_c = 1; + int32 temperature_f = 2; +} + +message LightSensorProto { + int32 link_quality = 1; + int32 luminosity = 2; +} + +message ClimateSensorProto { + int32 temperature_c = 1; + int32 humidity = 2; + int32 co2_level = 3; +} + +message SwitchSensorProto { + bool state = 1; +} diff --git a/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/services/collector_controller.proto b/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/services/collector_controller.proto index e69de29..a762e65 100644 --- a/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/services/collector_controller.proto +++ b/telemetry/serialization/proto-schemas/src/main/protobuf/telemetry/services/collector_controller.proto @@ -0,0 +1,18 @@ +syntax = "proto3"; + +package telemetry.service.collector; + +option java_multiple_files = true; +option java_package = "ru.yandex.practicum.grpc.telemetry.collector"; + +import "telemetry/messages/sensor_event.proto"; +import "telemetry/messages/hub_event.proto"; + +import "google/protobuf/empty.proto"; + +service CollectorController { + rpc CollectSensorEvent(telemetry.message.event.SensorEventProto) + returns (google.protobuf.Empty); + rpc CollectHubEvent(telemetry.message.event.HubEventProto) + returns (google.protobuf.Empty); +}