Files

171 lines
10 KiB
JavaScript
Raw Permalink 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 Producer — Эмулятор IoT-датчиков «Умного дома»
// =============================================================================
// Генерирует случайные события от 8 виртуальных датчиков (температура,
// влажность, энергопотребление) и отправляет их в очередь RabbitMQ.
//
// Архитектура пайплайна:
// Producer → RabbitMQ (iot-events) → Consumer → Redis + MongoDB → Dashboard
//
// Переменные окружения (задаются Terraform'ом через json_env):
// RABBIT_HOST — хост RabbitMQ (K8s service.cluster.local)
// RABBIT_PORT — порт AMQP (по умолчанию 5672)
// RABBIT_USER — логин RabbitMQ
// RABBIT_PASS — пароль RabbitMQ
// RABBIT_QUEUE — имя очереди (по умолчанию "iot-events")
// PRODUCE_INTERVAL — интервал между событиями в секундах (по умолчанию 3)
// PORT — порт HTTP-сервера (по умолчанию 3000)
//
// Эндпоинты:
// GET / — статистика (events_sent, errors)
// GET /health — health-check для Kubernetes
// =============================================================================
// ---------------------------------------------------------------------------
// Зависимости
// ---------------------------------------------------------------------------
const amqp = require("amqplib"); // AMQP-клиент для RabbitMQ (0-10-0)
const express = require("express"); // HTTP-фреймворк для /health и /stats
// ---------------------------------------------------------------------------
// Конфигурация RabbitMQ из переменных окружения
// Fallback-значения — для локальной разработки (localhost + guest:guest)
// ---------------------------------------------------------------------------
const HOST = process.env.RABBIT_HOST || "localhost"; // K8s-хост RabbitMQ
const PORT = parseInt(process.env.RABBIT_PORT || "5672"); // порт AMQP (не HTTP!)
const USER = process.env.RABBIT_USER || "guest"; // логин (дефолтный в RabbitMQ)
const PASS = process.env.RABBIT_PASS || "guest"; // пароль
const QUEUE = process.env.RABBIT_QUEUE || "iot-events"; // имя очереди
// Интервал генерации событий: из переменной окружения (в секундах) → миллисекунды
const INTERVAL = parseInt(process.env.PRODUCE_INTERVAL || "3") * 1000;
// ---------------------------------------------------------------------------
// Конфигурация виртуальных датчиков
// ---------------------------------------------------------------------------
// 8 датчиков, 3 типа:
// - temperature (3 шт.): жилая комната, спальня, офис
// - humidity (2 шт.): жилая комната, спальня
// - power_meter (3 шт.): кухня, жилая комната, офис
// Каждый датчик имеет диапазон min/max — в этих пределах генерируются значения
// ---------------------------------------------------------------------------
const SENSORS = [
// --- Датчики температуры (°C) ---
{ id: "temp_living", type: "temperature", location: "living_room", unit: "°C", min: 18, max: 28 },
{ id: "temp_bedroom", type: "temperature", location: "bedroom", unit: "°C", min: 16, max: 26 },
{ id: "temp_office", type: "temperature", location: "office", unit: "°C", min: 19, max: 27 },
// --- Датчики влажности (%) ---
{ id: "humidity_living", type: "humidity", location: "living_room", unit: "%", min: 30, max: 70 },
{ id: "humidity_bed", type: "humidity", location: "bedroom", unit: "%", min: 35, max: 65 },
// --- Датчики энергопотребления (kW) ---
{ id: "power_kitchen", type: "power_meter", location: "kitchen", unit: "kW", min: 0.1, max: 3.5 },
{ id: "power_living", type: "power_meter", location: "living_room", unit: "kW", min: 0.2, max: 2 },
{ id: "power_office", type: "power_meter", location: "office", unit: "kW", min: 0.3, max: 4 },
];
// ---------------------------------------------------------------------------
// Глобальное состояние
// ---------------------------------------------------------------------------
let channel = null; // AMQP-канал RabbitMQ (ленивая инициализация)
let eventsSent = 0; // счётчик отправленных событий (для /)
let errors = 0; // счётчик ошибок (для /)
// ---------------------------------------------------------------------------
// connect() — ленивое подключение к RabbitMQ
// ---------------------------------------------------------------------------
// Вызывается при старте и при обрыве соединения.
// Использует протокол amqp:// (НЕ amqps:// — трафик внутри K8s-кластера,
// TLS не нужен).
// assertQueue с durable: true — очередь переживёт рестарт RabbitMQ.
// ---------------------------------------------------------------------------
async function connect() {
// Формируем строку подключения: amqp://логин:пароль@хост:порт
const conn = await amqp.connect(`amqp://${USER}:${PASS}@${HOST}:${PORT}`);
// Создаём канал (виртуальное соединение внутри TCP-соединения)
channel = await conn.createChannel();
// Гарантируем существование очереди (идемпотентно)
// durable: true — сообщения не потеряются при рестарте брокера
await channel.assertQueue(QUEUE, { durable: true });
console.log(`Producer connected to RabbitMQ, queue: ${QUEUE}`);
}
// ---------------------------------------------------------------------------
// produce() — бесконечный цикл генерации и отправки событий
// ---------------------------------------------------------------------------
// Каждые INTERVAL миллисекунд:
// 1. Выбирает случайный датчик из SENSORS
// 2. Генерирует случайное значение в диапазоне [min, max]
// 3. Формирует JSON-событие с полями event_id, sensor_id, device_type,
// location, value, unit, timestamp
// 4. Отправляет в очередь RabbitMQ (persistent: true — не теряется
// при падении брокера)
// 5. Каждые 10 событий логирует прогресс
//
// При ошибке: инкрементирует счётчик errors и продолжает цикл.
// Канал пересоздаётся при обрыве (connect() вызывается заново).
// ---------------------------------------------------------------------------
async function produce() {
while (true) { // бесконечный цикл эмуляции
try {
// Если канала нет (первый запуск или обрыв) — переподключаемся
if (!channel) await connect();
// Выбираем случайный датчик
const s = SENSORS[Math.floor(Math.random() * SENSORS.length)];
// Формируем событие
const event = {
event_id: crypto.randomUUID(), // уникальный UUID для каждого события
sensor_id: s.id, // идентификатор датчика
device_type: s.type, // тип: temperature | humidity | power_meter
location: s.location, // локация: living_room | bedroom | office | kitchen
value: +(s.min + Math.random() * (s.max - s.min)).toFixed(2), // случайное значение, 2 знака
unit: s.unit, // единица измерения: °C | % | kW
timestamp: new Date().toISOString(), // ISO 8601 timestamp
};
// Отправляем в очередь RabbitMQ
// Buffer.from — сериализуем JSON в бинарный буфер
// persistent: true — сообщение будет записано на диск брокера
channel.sendToQueue(QUEUE, Buffer.from(JSON.stringify(event)), { persistent: true });
eventsSent++; // инкрементируем счётчик
// Каждые 10 событий — логируем (чтобы не засорять логи)
if (eventsSent % 10 === 0) console.log(`Sent ${eventsSent} events`);
} catch (e) {
// При ошибке (обрыв соединения с RabbitMQ, etc.)
errors++;
console.error("Error:", e.message);
// Канал сбрасываем — в следующей итерации connect() переподключит
}
// Ждём INTERVAL миллисекунд перед следующей итерацией
// НЕ используем setInterval — чтобы избежать наложения итераций
await new Promise(r => setTimeout(r, INTERVAL));
}
}
// ---------------------------------------------------------------------------
// HTTP-сервер (Express)
// ---------------------------------------------------------------------------
const app = express();
// GET / — статистика producer'а (используется для мониторинга)
app.get("/", (req, res) => res.json({
service: "iot-producer",
status: "running",
queue: QUEUE,
events_sent: eventsSent,
errors
}));
// GET /health — health-check для Kubernetes (liveness/readiness probe)
// Возвращает 200 если процесс жив, без проверки RabbitMQ
app.get("/health", (req, res) => res.json({ status: "ok" }));
// Запуск сервера: слушаем порт из env или 3000, затем стартуем эмуляцию
app.listen(process.env.PORT || 3000, () => {
console.log("Producer started");
produce(); // запускаем бесконечный цикл генерации событий
});