# IoT Consumer — обработчик событий Читает сообщения из очереди RabbitMQ и сохраняет их в Redis (оперативный кэш) и MongoDB (постоянный архив). ## Пайплайн ``` Producer → RabbitMQ → Consumer (этот сервис) → Redis + MongoDB → Dashboard ``` ## Хранилища ### Redis — оперативный кэш | Ключ | Тип | Содержимое | |---|---|---| | `iot:latest` | Hash | Последнее значение каждого датчика (`sensor_id → JSON`) | | `iot:counters` | Hash | Счётчики событий по типам (`temperature: N`, ...) | | `iot:recent` | ZSET | Лента последних 1000 событий (сортировка по времени) | ### MongoDB — постоянный архив - Коллекция: `iot_events` - TTL-индекс: автоматическое удаление через **7 дней** ## Переменные окружения | Переменная | По умолчанию | Описание | |---|---|---| | `RMQ_HOST` | `localhost` | Хост RabbitMQ | | `RMQ_PORT` | `5672` | Порт AMQP | | `RMQ_USER` | `guest` | Логин RabbitMQ | | `RMQ_PASS` | `guest` | Пароль RabbitMQ | | `RMQ_QUEUE` | `iot-events` | Имя очереди | | `REDIS_HOST` | `localhost` | Хост Redis | | `REDIS_PORT` | `6379` | Порт Redis | | `REDIS_PASS` | — | Пароль Redis (необязательно) | | `MONGO_URI` | `mongodb://localhost:27017/iot` | Строка подключения MongoDB | | `MONGO_DB` | `iot` | Имя базы MongoDB | | `PORT` | `3000` | Порт HTTP | ## Эндпоинты | Путь | Ответ | |---|---| | `GET /` | `{"service":"iot-consumer","consumed":N,"errors":N}` | | `GET /health` | `{"status":"ok"}` | ## Локальный запуск ```bash npm install # Нужны RabbitMQ, Redis и MongoDB (локально или в Docker) RMQ_HOST=localhost REDIS_HOST=localhost MONGO_URI=mongodb://localhost:27017/iot node server.js ``` ## Деплой на Nubes Cloud 1. Запушьте этот репо в свой Gitea 2. В Terraform-конфиге укажите `git_path` на ваш репо и `git_revision` 3. Nubes автоматически запустит `npm start` ## HTTP-эндпоинт Сейчас `GET /` отдаёт только статистику: `{"consumed": N, "errors": N}`. Nubes Cloud автоматически создаёт публичный домен для сервиса. При необходимости можно расширить: - `/metrics` — Prometheus-формат для Grafana - `/health` — проверка доступности RabbitMQ + Redis + MongoDB - SSE-стриминг обработанных событий в реальном времени ## Отказоустойчивость - Каждое хранилище (Redis, MongoDB) обёрнуто в отдельный `try/catch` - Отказ одного хранилища **не ломает** запись в другое - Ошибки логируются, consumer продолжает работу