// ============================================================================= // 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(); // запускаем бесконечный цикл генерации событий });