Files
tf-iot-consumer/server.js
T

269 lines
14 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// =============================================================================
// IoT Consumer — Обработчик событий из RabbitMQ
// =============================================================================
// Читает сообщения из очереди RabbitMQ (iot-events), помещённые туда
// Producer'ом, и сохраняет их в два хранилища:
// 1. Redis — оперативный кэш для Dashboard:
// - iot:latest (hash) — последнее значение каждого датчика
// - iot:counters (hash) — счётчики событий по типам устройств
// - iot:recent (zset) — лента последних 1000 событий
// 2. MongoDB — постоянный архив всех событий:
// - коллекция iot_events с TTL-индексом (автоудаление через 7 дней)
//
// Каждое хранилище обёрнуто в try/catch — отказ одного не ломает другое.
//
// Переменные окружения (задаются Terraform'ом через json_env):
// RMQ_HOST — хост RabbitMQ
// RMQ_PORT — порт AMQP (5672)
// RMQ_USER — логин RabbitMQ
// RMQ_PASS — пароль RabbitMQ
// RMQ_QUEUE — имя очереди (iot-events)
// REDIS_HOST — хост Redis
// REDIS_PORT — порт Redis (6379)
// REDIS_PASS — пароль Redis
// MONGO_URI — строка подключения MongoDB (mongodb://user:pass@host:port/db)
// MONGO_DB — имя базы MongoDB (iot)
// PORT — порт HTTP-сервера (3000)
//
// Эндпоинты:
// GET / — статистика (consumed, errors)
// GET /health — health-check для Kubernetes
// =============================================================================
// ---------------------------------------------------------------------------
// Зависимости
// ---------------------------------------------------------------------------
const amqp = require("amqplib"); // AMQP-клиент для RabbitMQ
const express = require("express"); // HTTP-фреймворк для /health и /stats
const { createClient } = require("redis"); // Redis-клиент (v4+, с Promise API)
const { MongoClient } = require("mongodb");// MongoDB-клиент (v6+, нативный драйвер)
// ---------------------------------------------------------------------------
// Конфигурация RabbitMQ из переменных окружения
// ---------------------------------------------------------------------------
const RMQ_HOST = process.env.RMQ_HOST || "localhost"; // K8s-хост
const RMQ_PORT = parseInt(process.env.RMQ_PORT || "5672"); // AMQP-порт
const RMQ_USER = process.env.RMQ_USER || "guest"; // логин
const RMQ_PASS = process.env.RMQ_PASS || "guest"; // пароль
const RMQ_QUEUE = process.env.RMQ_QUEUE || "iot-events"; // имя очереди
// ---------------------------------------------------------------------------
// Конфигурация Redis из переменных окружения
// ---------------------------------------------------------------------------
const REDIS_HOST = process.env.REDIS_HOST || "localhost"; // K8s-хост
const REDIS_PORT = parseInt(process.env.REDIS_PORT || "6379"); // порт Redis
const REDIS_PASS = process.env.REDIS_PASS || undefined; // пароль (undefined = без пароля)
// ---------------------------------------------------------------------------
// Конфигурация MongoDB из переменных окружения
// ---------------------------------------------------------------------------
// MONGO_URI формируется Terraform'ом как:
// mongodb://admin:@хост:27017/iot?authSource=admin
// (без пароля для демо-стенда — platform limitation)
// ---------------------------------------------------------------------------
const MONGO_URI = process.env.MONGO_URI || "mongodb://localhost:27017/iot";
const MONGO_DB = process.env.MONGO_DB || "iot";
// ---------------------------------------------------------------------------
// Глобальное состояние
// ---------------------------------------------------------------------------
let consumed = 0, errors = 0; // счётчики для /
let redis, // Redis-клиент (ленивая инициализация)
mongo, // MongoDB-клиент (ленивая инициализация)
mongoColl; // MongoDB-коллекция iot_events
// ---------------------------------------------------------------------------
// getRedis() — ленивое подключение к Redis
// ---------------------------------------------------------------------------
// Создаёт клиент Redis один раз, при первом вызове.
// Подавляем ошибки Redis (redis.on("error", () => {})) — не хотим крашить
// consumer при временной недоступности Redis.
// ---------------------------------------------------------------------------
async function getRedis() {
if (!redis) {
// Создаём клиент с настройками из переменных окружения
redis = createClient({
socket: { host: REDIS_HOST, port: REDIS_PORT }, // TCP-сокет
password: REDIS_PASS, // пароль (undefined если нет)
});
// Подавляем ошибки Redis — не крашим consumer из-за временных проблем
redis.on("error", () => {});
await redis.connect();
console.log("Redis connected");
}
return redis;
}
// ---------------------------------------------------------------------------
// getMongo() — ленивое подключение к MongoDB
// ---------------------------------------------------------------------------
// Создаёт клиент MongoDB один раз, при первом вызове.
// При подключении создаёт TTL-индекс на поле timestamp:
// expireAfterSeconds: 604800 = 7 дней
// Это значит MongoDB автоматически удалит документы старше 7 дней —
// не нужно ручной очистки.
// ---------------------------------------------------------------------------
async function getMongo() {
if (!mongo) {
mongo = new MongoClient(MONGO_URI); // строка подключения из env
await mongo.connect();
// Получаем коллекцию iot_events в базе MONGO_DB
mongoColl = mongo.db(MONGO_DB).collection("iot_events");
// TTL-индекс: документы старше 7 дней удаляются автоматически
await mongoColl.createIndex({ timestamp: 1 }, { expireAfterSeconds: 604800 });
console.log("MongoDB connected");
}
return mongoColl;
}
// ---------------------------------------------------------------------------
// store(event) — сохраняет событие в Redis и MongoDB
// ---------------------------------------------------------------------------
// Каждое хранилище — в своём try/catch:
// - Если Redis упал — событие всё равно попадёт в MongoDB
// - Если MongoDB упала — событие всё равно попадёт в Redis
// Это гарантирует, что отказ одного хранилища не блокирует другое.
//
// Redis-структуры:
// iot:latest — HSET sensor_id → JSON{value, unit, location, timestamp}
// (каждый новый запрос перезаписывает предыдущий)
// iot:counters — HINCRBY device_type → +1
// (счётчик событий по типам: temperature, humidity, power_meter)
// iot:recent — ZADD score=Date.now() → JSON-событие
// ZREMRANGEBYRANK 0..-1001 — держим только 1000 последних
// ---------------------------------------------------------------------------
async function store(event) {
// === Запись в Redis (оперативный кэш) ===
try {
const r = await getRedis();
// HSET — последнее значение датчика (ключ = sensor_id)
// Перезаписывает предыдущее значение — всегда актуальное
await r.hSet("iot:latest", event.sensor_id, JSON.stringify({
value: event.value,
unit: event.unit,
location: event.location,
timestamp: event.timestamp,
}));
// HINCRBY — атомарный инкремент счётчика по типу устройства
await r.hIncrBy("iot:counters", event.device_type, 1);
// ZADD — добавляем в сортированное множество (счёт = timestamp)
// score = Date.now() для сортировки по времени
await r.zAdd("iot:recent", [{
score: Date.now(),
value: JSON.stringify(event),
}]);
// ZREMRANGEBYRANK — удаляем всё что выходит за пределы последних 1000
// 0..-1001 означает: удалить с 0-го элемента до элемента с индексом
// (len - 1001), оставляя 1000 последних
await r.zRemRangeByRank("iot:recent", 0, -1001);
} catch (e) {
// Redis недоступен — логируем, но НЕ крашим consumer
console.error("Redis error:", e.message);
}
// === Запись в MongoDB (постоянный архив) ===
try {
const coll = await getMongo();
// insertOne — вставка одного документа
// timestamp преобразуется из ISO-строки в Date для TTL-индекса
await coll.insertOne({
event_id: event.event_id,
sensor_id: event.sensor_id,
device_type: event.device_type,
location: event.location,
value: event.value,
unit: event.unit,
timestamp: new Date(event.timestamp), // Date-объект для TTL
});
} catch (e) {
// MongoDB недоступна — логируем, но НЕ крашим consumer
console.error("Mongo error:", e.message);
}
}
// ---------------------------------------------------------------------------
// start() — подключение к RabbitMQ и запуск консьюмера
// ---------------------------------------------------------------------------
// 1. Подключается к RabbitMQ по AMQP
// 2. Создаёт канал и гарантирует существование очереди
// 3. Устанавливает prefetch(10) — не более 10 неподтверждённых сообщений
// одновременно (backpressure: не загружаем consumer слишком сильно)
// 4. consume() — подписывается на очередь, обрабатывает каждое сообщение:
// a. Парсит JSON из тела сообщения
// b. Вызывает store() для записи в Redis + MongoDB
// c. ack() — подтверждает обработку (сообщение удаляется из очереди)
// d. При ошибке: nack(msg, false, true) — возвращает в очередь (requeue)
// ---------------------------------------------------------------------------
async function start() {
// Подключение к RabbitMQ
const conn = await amqp.connect(
`amqp://${RMQ_USER}:${RMQ_PASS}@${RMQ_HOST}:${RMQ_PORT}`
);
const ch = await conn.createChannel();
// Гарантируем существование очереди (идемпотентно)
await ch.assertQueue(RMQ_QUEUE, { durable: true });
// prefetch(10): не отправлять больше 10 сообщений одновременно,
// пока предыдущие не подтверждены (ack)
// Это защита от перегрузки consumer'а
ch.prefetch(10);
// Подписка на очередь
await ch.consume(RMQ_QUEUE, async (msg) => {
// msg === null означает отмену подписки (например, очередь удалена)
if (!msg) return;
try {
// Парсим JSON из тела сообщения (Buffer → string → object)
const event = JSON.parse(msg.content.toString());
// Сохраняем в Redis + MongoDB
await store(event);
consumed++; // инкрементируем счётчик успешных
ch.ack(msg); // подтверждаем — сообщение удаляется из очереди
// Каждые 10 сообщений — логируем прогресс
if (consumed % 10 === 0) console.log(`Consumed ${consumed} events`);
} catch (e) {
errors++;
// nack с requeue=true — возвращаем сообщение в очередь
// для повторной попытки (если ошибка временная)
// ВНИМАНИЕ: при систематической ошибке может вызвать infinite loop!
ch.nack(msg, false, true);
console.error("Error:", e.message);
}
});
console.log(`Consumer listening on queue: ${RMQ_QUEUE}`);
}
// ---------------------------------------------------------------------------
// HTTP-сервер (Express)
// ---------------------------------------------------------------------------
const app = express();
// GET / — статистика consumer'а
app.get("/", (req, res) => res.json({
service: "iot-consumer",
status: "running",
queue: RMQ_QUEUE,
consumed,
errors,
}));
// GET /health — для Kubernetes liveness/readiness probe
app.get("/health", (req, res) => res.json({ status: "ok" }));
// Запуск: слушаем порт, затем подключаемся к RabbitMQ
app.listen(process.env.PORT || 3000, () => {
console.log("Consumer started");
start(); // подключаемся к RabbitMQ и начинаем обрабатывать сообщения
});