first commit
This commit is contained in:
commit
628642a1f7
|
|
@ -0,0 +1,19 @@
|
|||
# Environment variables — JANGAN pernah di-commit!
|
||||
.env
|
||||
.env.local
|
||||
.env.*.local
|
||||
|
||||
# Dependencies
|
||||
node_modules/
|
||||
|
||||
# Logs
|
||||
*.log
|
||||
npm-debug.log*
|
||||
|
||||
# OS files
|
||||
.DS_Store
|
||||
Thumbs.db
|
||||
|
||||
# Temporary test/scratch files
|
||||
test-latency-standalone.js
|
||||
scratch-check.js
|
||||
|
|
@ -0,0 +1,626 @@
|
|||
# 🍄 Backend IoT Kumbung Jamur
|
||||
|
||||
Backend Node.js untuk sistem monitoring dan kontrol otomatis kumbung jamur berbasis IoT. Sistem ini menghubungkan sensor **SHT31** dan relay pada **ESP32** ke aplikasi mobile **Flutter** melalui **MQTT (HiveMQ Cloud)**, dengan penjadwalan penyiraman otomatis menggunakan **BullMQ (Redis)**, sistem push notification anti-spam via **OneSignal**, serta database utama **Supabase**.
|
||||
|
||||
Aplikasi ini dirancang dengan arsitektur **High Availability (Dual Server Failover)** dan **In-Memory Optimizations** untuk memastikan keandalan sistem yang optimal dan efisiensi resource yang sangat tinggi.
|
||||
|
||||
---
|
||||
|
||||
## 🏗️ Arsitektur Sistem Terdistribusi
|
||||
|
||||
Sistem ini mendukung arsitektur **Failover Otomatis** menggunakan dua server (Primary & Backup) untuk menjamin layanan tetap berjalan meskipun server utama mengalami kendala (*down*).
|
||||
|
||||
```
|
||||
┌──────────────────────────────────────┐
|
||||
│ Aplikasi Flutter │
|
||||
└──────┬────────────────────────┬──────┘
|
||||
│
|
||||
HTTP REST API │
|
||||
(https://...railway.app)
|
||||
▼
|
||||
┌────────────────────────┐
|
||||
│ PRIMARY SERVER ├ ────────────────────
|
||||
│ (Railway) │ |
|
||||
└───────────┬────────────┘ |
|
||||
│ │
|
||||
┌─────┴─────────────────────────────────┴─────┐
|
||||
│ │
|
||||
▼ ▼ ▼
|
||||
Supabase (DB) Upstash (Redis) HiveMQ Cloud
|
||||
┌────────────────┐ ┌─────────────────┐ ┌──────────────────┐
|
||||
│ Tabel Utama & │ │ Antrian Kerja │ │ Broker MQTT │
|
||||
│ RPC Functions │ │ (BullMQ) │ │ (TLS) │
|
||||
└────────────────┘ └─────────────────┘ └────────┬─────────┘
|
||||
│
|
||||
▼ (mqtts://)
|
||||
┌─────────────────┐
|
||||
│ ESP32 + SHT31 │
|
||||
└─────────────────┘
|
||||
```
|
||||
|
||||
### 🔄 Alur & Fitur Utama Failover (Dynamic Failover Manager)
|
||||
1. **Primary Server (Railway)**: Berjalan secara aktif. Menangani REST API, koneksi MQTT, pemrosesan antrian BullMQ, dan pendeteksian perangkat offline.
|
||||
2. **Backup Server (VPS - Standby)**:
|
||||
- Saat startup, background services (MQTT, Worker BullMQ, Scheduler Restore, Offline Detector) **dinonaktifkan** (standby).
|
||||
- Secara berkala (setiap 10 detik / `FAILOVER_PING_INTERVAL_MS`), server backup mengirimkan ping (HTTP GET) ke `PRIMARY_SERVER_URL`.
|
||||
- **Jika Primary Down (Ping Gagal)**: Backup server langsung mengaktifkan semua background services lokal, me-resume worker BullMQ, me-restore jadwal aktif dari database, dan mengambil alih kendali sistem.
|
||||
- **Jika Primary Kembali Online (Ping Sukses)**: Backup server secara otomatis mendegradasi diri kembali ke mode standby (disconnect MQTT, pause worker, matikan offline detector) untuk menghindari bentrokan eksekusi (*race conditions*).
|
||||
|
||||
---
|
||||
|
||||
## ⚡ Optimalisasi Kinerja & Efisiensi Database
|
||||
|
||||
Untuk mencegah pembengkakan biaya database (kuota request Supabase) dan beban broker MQTT akibat frekuensi transmisi data ESP32 yang tinggi, backend mengimplementasikan beberapa teknik optimalisasi tingkat tinggi:
|
||||
|
||||
### 1. In-Memory Cache (Threshold Cache)
|
||||
* **Masalah**: ESP32 mengirim data sensor secara real-time. Jika backend harus menanyakan batas threshold ke Supabase setiap kali data masuk, database akan terbebani ribuan query per menit.
|
||||
* **Solusi**: Nilai threshold disimpan dalam *in-memory cache* server dengan TTL 30 detik (`CACHE_TTL_MS`). Query ke Supabase hanya dilakukan saat cache kedaluwarsa atau terjadi update threshold melalui API (cache langsung di-invalidate secara instan).
|
||||
|
||||
### 2. Smart Filtering & Deadband Logging (Sensor Log Cache)
|
||||
Backend menyaring log sensor sebelum disimpan ke database Supabase melalui algoritma *Deadband filtering*:
|
||||
* Data sensor hanya akan disimpan ke tabel `sensor_logs` jika memenuhi salah satu kondisi berikut:
|
||||
- Suhu berubah lebih dari **0.5°C** (`TEMP_DELTA`).
|
||||
- Kelembapan berubah lebih dari **1.0%** (`HUM_DELTA`).
|
||||
- Status relay berubah (`ON` ↔ `OFF`).
|
||||
- Mode operasi perangkat berubah (`auto`, `manual`, `offline`).
|
||||
- Interval waktu detak jantung (*heartbeat*) telah mencapai **5 menit** (`HEARTBEAT_INTERVAL_MS`), bertujuan untuk menjaga kontinuitas grafik di UI.
|
||||
* Mengurangi penyimpanan database hingga **90%** tanpa kehilangan data historis yang penting.
|
||||
|
||||
### 3. Online Status Throttling
|
||||
* Status keaktifan perangkat (`last_seen` dan `is_online`) di database diperbarui secara berkala maksimal **1 menit sekali** (`LAST_SEEN_INTERVAL_MS`), mencegah spam penulisan (*write-heavy*) ke Supabase.
|
||||
|
||||
---
|
||||
|
||||
## 🚨 Pendeteksi Perangkat Offline & Push Notification
|
||||
|
||||
### 1. Offline Detector Job
|
||||
* Berjalan di latar belakang setiap 5 menit.
|
||||
* Memindai perangkat di database yang berstatus `is_online = true` namun `last_seen` berumur lebih dari 5 menit lalu.
|
||||
* Jika ditemukan, secara otomatis memperbarui status perangkat menjadi offline di database dan memicu push notification darurat ke pemilik perangkat.
|
||||
|
||||
### 2. Anti-Spam Push Notification (OneSignal Integration)
|
||||
* **Koneksi**: Integrasi langsung ke OneSignal menggunakan API REST resmi.
|
||||
* **Anti-Spam Guard (Cooldown)**: Setiap notifikasi yang sama (misal peringatan sensor kering/panas atau status offline) memiliki *cooldown time* (seperti 30 menit atau 1 jam) yang dikelola di Redis / In-Memory Map agar pengguna tidak dibanjiri spam push notification.
|
||||
* **Notification Stacking & Threading**: Menggunakan `android_group`, `thread_id`, dan `collapse_id` dengan nama `jamur_monitoring_group` untuk mengelompokkan notifikasi secara rapi di notification tray Android & iOS (seperti gaya chat WhatsApp).
|
||||
|
||||
---
|
||||
|
||||
## 📋 Prasyarat Layanan Cloud
|
||||
|
||||
Pastikan Anda memiliki akun dan konfigurasi untuk layanan-layanan berikut:
|
||||
|
||||
| Layanan | Keterangan |
|
||||
|---|---|
|
||||
| [Supabase](https://supabase.com) | Database PostgreSQL utama (gratis) |
|
||||
| [HiveMQ Cloud](https://www.hivemq.com/mqtt-cloud-broker/) | MQTT Broker TLS aman port 8883 (gratis) |
|
||||
| [Upstash Redis](https://upstash.com) | Redis Cloud untuk antrian BullMQ (gratis) |
|
||||
| [OneSignal](https://onesignal.com) | Platform Push Notification ke Aplikasi Mobile |
|
||||
| Node.js >= 18 | Runtime JavaScript lokal |
|
||||
|
||||
---
|
||||
|
||||
## ⚙️ Instalasi & Konfigurasi
|
||||
|
||||
### 1. Clone & Install Dependencies
|
||||
|
||||
```bash
|
||||
git clone <repository-url>
|
||||
cd backend-jamur
|
||||
npm install
|
||||
```
|
||||
|
||||
### 2. Setup Environment Variables
|
||||
|
||||
Buat file `.env` di root folder aplikasi, lalu isi konfigurasi berikut:
|
||||
|
||||
```env
|
||||
# 💻 Server Configuration
|
||||
PORT=3000
|
||||
IS_BACKUP_SERVER=false # Set 'true' jika ini dideploy ke server VPS Backup
|
||||
PRIMARY_SERVER_URL=https://nama-project-kamu.up.railway.app # URL Primary Server (diperlukan jika ini Backup Server)
|
||||
FAILOVER_PING_INTERVAL_MS=10000 # Interval cek primary server (dalam milidetik, default 10 detik)
|
||||
|
||||
# ⚡ Supabase Configuration (Ambil dari: Project Settings > API)
|
||||
SUPABASE_URL=https://xxxxxxxxxxxxxxxx.supabase.co
|
||||
SUPABASE_SERVICE_KEY=sb_secret_xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx
|
||||
|
||||
# 📡 HiveMQ Cloud MQTT Broker (Ambil dari: Clusters > MQTT Credentials & Connection Settings)
|
||||
MQTT_BROKER_URL=mqtts://xxxxxxxxxxxxxxxxxxxxxxxx.s1.eu.hivemq.cloud:8883
|
||||
MQTT_PORT=8883
|
||||
MQTT_USERNAME=username_hivemq_kamu
|
||||
MQTT_PASSWORD=password_hivemq_kamu
|
||||
|
||||
# 🔴 Upstash Redis Connection (Ambil dari: Database > Details > Connection URL)
|
||||
REDIS_URL=rediss://default:xxxxxxxx@xxxx.upstash.io:6379
|
||||
|
||||
# 🔔 OneSignal Push Notification (Ambil dari: Settings > Keys & IDs)
|
||||
ONESIGNAL_APP_ID=xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx
|
||||
ONESIGNAL_REST_API_KEY=xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx
|
||||
```
|
||||
|
||||
> ⚠️ **PENTING**: Jangan pernah melakukan commit file `.env` ke Git! File ini sudah otomatis diabaikan di file `.gitignore`.
|
||||
|
||||
### 3. Setup Database Supabase (Skema Tabel & Stored Procedure)
|
||||
|
||||
Jalankan perintah SQL berikut di dashboard **Supabase > SQL Editor**:
|
||||
|
||||
```sql
|
||||
-- 1. TABEL UTAMA: Perangkat ESP32
|
||||
CREATE TABLE devices (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
device_id TEXT UNIQUE NOT NULL, -- e.g. "esp32-01"
|
||||
label TEXT, -- nama display e.g. "Kumbung Barat"
|
||||
location TEXT, -- lokasi penempatan
|
||||
claim_code TEXT UNIQUE, -- kode klaim huruf kapital, e.g. "JAMUR01"
|
||||
claimed_by UUID REFERENCES auth.users(id),
|
||||
claimed_at TIMESTAMPTZ,
|
||||
is_online BOOLEAN DEFAULT false,
|
||||
last_seen TIMESTAMPTZ,
|
||||
current_mode TEXT DEFAULT 'auto' -- mode kerja aktif: auto, manual, offline
|
||||
);
|
||||
|
||||
-- 2. TABEL LOG: Riwayat Sensor & Log Aksi
|
||||
CREATE TABLE sensor_logs (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
device_id TEXT NOT NULL,
|
||||
temperature FLOAT,
|
||||
humidity FLOAT,
|
||||
relay_state BOOLEAN DEFAULT false, -- true = ON, false = OFF
|
||||
mode TEXT DEFAULT 'auto', -- snapshot mode kerja saat log terekam
|
||||
event TEXT, -- event penting (e.g. "manual_stop", "system_on")
|
||||
note TEXT, -- catatan tambahan
|
||||
created_at TIMESTAMPTZ DEFAULT NOW()
|
||||
);
|
||||
|
||||
-- 3. TABEL CONFIG: Batas Sensor Otomatisasi
|
||||
CREATE TABLE thresholds (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
device_id TEXT UNIQUE NOT NULL,
|
||||
temp_max FLOAT NOT NULL DEFAULT 30.0, -- suhu maks sebelum pompa ON
|
||||
hum_max FLOAT NOT NULL DEFAULT 80.0, -- kelembapan maks sebelum pompa ON
|
||||
updated_at TIMESTAMPTZ DEFAULT NOW()
|
||||
);
|
||||
|
||||
-- 4. TABEL JADWAL: Jadwal Penyiraman BullMQ
|
||||
CREATE TABLE schedules (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
device_id TEXT NOT NULL,
|
||||
label TEXT,
|
||||
cron TEXT NOT NULL, -- cron format: "menit jam hari bulan hari-minggu"
|
||||
duration_s INTEGER NOT NULL, -- durasi penyiraman dalam detik
|
||||
bull_job_id TEXT, -- ID job BullMQ untuk kontrol sinkronisasi
|
||||
is_active BOOLEAN DEFAULT true,
|
||||
created_at TIMESTAMPTZ DEFAULT NOW()
|
||||
);
|
||||
|
||||
-- 5. FUNCTION: Rata-rata Harian (RPC Function)
|
||||
CREATE OR REPLACE FUNCTION get_daily_average(p_device_id TEXT, p_days INT)
|
||||
RETURNS TABLE(day DATE, avg_temp NUMERIC, avg_hum NUMERIC) AS $$
|
||||
BEGIN
|
||||
RETURN QUERY
|
||||
SELECT
|
||||
created_at::DATE AS day,
|
||||
ROUND(AVG(temperature)::NUMERIC, 1) AS avg_temp,
|
||||
ROUND(AVG(humidity)::NUMERIC, 1) AS avg_hum
|
||||
FROM sensor_logs
|
||||
WHERE device_id = p_device_id
|
||||
AND created_at >= NOW() - (p_days || ' days')::INTERVAL
|
||||
AND temperature IS NOT NULL
|
||||
AND humidity IS NOT NULL
|
||||
GROUP BY created_at::DATE
|
||||
ORDER BY day DESC;
|
||||
END;
|
||||
$$ LANGUAGE plpgsql;
|
||||
|
||||
-- 6. FUNCTION: Rata-rata Per Jam (RPC Function)
|
||||
CREATE OR REPLACE FUNCTION get_hourly_average(p_device_id TEXT, p_days INT)
|
||||
RETURNS TABLE(hour TIMESTAMPTZ, avg_temp NUMERIC, avg_hum NUMERIC) AS $$
|
||||
BEGIN
|
||||
RETURN QUERY
|
||||
SELECT
|
||||
date_trunc('hour', created_at) AS hour,
|
||||
ROUND(AVG(temperature)::NUMERIC, 1) AS avg_temp,
|
||||
ROUND(AVG(humidity)::NUMERIC, 1) AS avg_hum
|
||||
FROM sensor_logs
|
||||
WHERE device_id = p_device_id
|
||||
AND created_at >= NOW() - (p_days || ' days')::INTERVAL
|
||||
AND temperature IS NOT NULL
|
||||
AND humidity IS NOT NULL
|
||||
GROUP BY date_trunc('hour', created_at)
|
||||
ORDER BY hour DESC;
|
||||
END;
|
||||
$$ LANGUAGE plpgsql;
|
||||
```
|
||||
|
||||
### 4. Jalankan Server Secara Lokal
|
||||
|
||||
```bash
|
||||
# Jalankan mode Development (Auto-restart via nodemon)
|
||||
npm run dev
|
||||
|
||||
# Jalankan mode Production
|
||||
npm start
|
||||
```
|
||||
|
||||
Server akan aktif pada port `http://localhost:3000` (atau sesuai konfigurasi env `PORT`).
|
||||
|
||||
---
|
||||
|
||||
## 🚀 Panduan Deployment
|
||||
|
||||
### A. Deploy ke Railway (Sebagai Primary Server)
|
||||
1. Hubungkan repository GitHub Anda ke [Railway.app](https://railway.app).
|
||||
2. Buat Project Baru dan pilih **Deploy dari GitHub**.
|
||||
3. Di tab **Variables**, masukkan semua Environment Variables seperti isi file `.env` di atas (kecuali `PORT` karena dikelola otomatis oleh Railway).
|
||||
4. Klik **Generate Domain** di tab **Settings > Networking** untuk mendapatkan URL server publik (contoh: `https://nama-project-kamu.up.railway.app`).
|
||||
|
||||
### B. Deploy ke VPS (Sebagai Backup Server)
|
||||
1. Siapkan server VPS (Ubuntu/Debian) dengan Node.js >= 18 dan PM2 terinstall.
|
||||
2. Clone repository, jalankan `npm install`.
|
||||
3. Set file `.env` dengan variabel `IS_BACKUP_SERVER=true` dan isikan `PRIMARY_SERVER_URL` dengan URL Railway yang didapatkan dari langkah di atas.
|
||||
4. Jalankan server menggunakan PM2 agar berjalan di latar belakang:
|
||||
```bash
|
||||
pm2 start src/index.js --name backend-jamur-backup
|
||||
pm2 save
|
||||
pm2 startup
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 🔌 API Reference
|
||||
|
||||
Base URL (Development): `http://localhost:3000/api`
|
||||
Base URL (Production): `https://nama-project-kamu.up.railway.app/api`
|
||||
|
||||
> 🛡️ **Rate Limiting**:
|
||||
> * Global rate limit untuk semua API: **100 request/menit per IP**.
|
||||
> * Rate limit ketat untuk trigger siram manual: **5 request/menit per IP** (menghindari banjir air pada kumbung).
|
||||
|
||||
---
|
||||
|
||||
### 📋 Ringkasan Daftar API (API Cheat Sheet)
|
||||
|
||||
Untuk mempermudah pencarian, berikut adalah ringkasan seluruh endpoint API yang tersedia pada backend ini:
|
||||
|
||||
| Kategori | Fitur / Kegunaan | Method | Endpoint | Parameter | Deskripsi Singkat |
|
||||
|---|---|:---:|---|---|---|
|
||||
| **📱 Device** | Klaim Device Baru (Pairing) | `POST` | `/api/device/claim` | `claim_code` (Body), `user_id` (Body) | Menghubungkan perangkat fisik ke akun user via kode klaim. |
|
||||
| | Ambil Info Device User | `GET` | `/api/device/my-device/:userId` | `userId` (Path) | Mengambil detail perangkat milik user (status online, last seen). |
|
||||
| | Cek Status Online Device | `GET` | `/api/device/status/:deviceId` | `deviceId` (Path) | Membaca status keaktifan perangkat secara real-time. |
|
||||
| **🌡️ Threshold** | Ambil Threshold Aktif | `GET` | `/api/threshold/:deviceId` | `deviceId` (Path) | Membaca batas suhu & kelembapan otomatisasi aktif. |
|
||||
| | Update Threshold | `POST` | `/api/threshold/:deviceId` | `deviceId` (Path), `temp_max` (Body), `hum_max` (Body), `temp_min` (Body, opsional) | Mengubah batas threshold (langsung sinkron ke ESP32 via MQTT). |
|
||||
| **🗓️ Jadwal** | Ambil Semua Jadwal | `GET` | `/api/schedule/:deviceId` | `deviceId` (Path) | Membaca seluruh daftar jadwal penyiraman perangkat. |
|
||||
| | Buat Jadwal Baru | `POST` | `/api/schedule/:deviceId` | `deviceId` (Path), `cron` (Body), `duration_s` (Body), `label` (Body, opsional) | Mendaftarkan jadwal berulang baru ke DB & BullMQ (cron). |
|
||||
| | Hapus Jadwal | `DELETE` | `/api/schedule/:id` | `id` (Path) | Menghapus jadwal permanen dari database dan antrian BullMQ. |
|
||||
| | Toggle Status Jadwal | `PATCH` | `/api/schedule/:id/toggle` | `id` (Path) | Mengaktifkan/menonaktifkan jadwal sementara tanpa menghapus. |
|
||||
| | Trigger Siram Instan | `POST` | `/api/schedule/:deviceId/now` | `deviceId` (Path), `duration_s` (Body, opsional) | Memicu penyiraman manual sekali jalan (default 30 detik). |
|
||||
| | Hentikan Pompa Paksa | `POST` | `/api/schedule/:deviceId/stop` | `deviceId` (Path) | Mengirim sinyal `OFF` langsung ke relay pompa via MQTT. |
|
||||
| **⚙️ Mode** | Ambil Mode Kerja Aktif | `GET` | `/api/mode/:deviceId` | `deviceId` (Path) | Membaca mode kerja aktif perangkat (`auto`, `manual`, `offline`). |
|
||||
| | Ubah Mode Kerja | `POST` | `/api/mode/:deviceId` | `deviceId` (Path), `mode` (Body) | Mengubah mode kerja ESP32 dan sinkronisasi perintah via MQTT. |
|
||||
| **📊 Riwayat** | Ambil Log Sensor Terbaru | `GET` | `/api/history/:deviceId` | `deviceId` (Path), `limit` (Query, opsional) | Mengambil data sensor real-time terbaru (terlimit maks 500). |
|
||||
| | Rata-rata Harian | `GET` | `/api/history/:deviceId/daily` | `deviceId` (Path), `days` (Query, opsional) | Mengambil data rata-rata suhu & kelembapan harian (tren grafik). |
|
||||
| | Rata-rata Per Jam | `GET` | `/api/history/:deviceId/hourly` | `deviceId` (Path), `days` (Query, opsional) | Mengambil data rata-rata suhu & kelembapan per jam (grafik analitis). |
|
||||
|
||||
---
|
||||
|
||||
### 📱 Device Management
|
||||
|
||||
#### 1. Klaim Device Baru (Pairing)
|
||||
Menghubungkan kode unik perangkat fisik dengan ID akun pengguna.
|
||||
* **Endpoint**: `POST /api/device/claim`
|
||||
* **Body Request**:
|
||||
```json
|
||||
{
|
||||
"claim_code": "JAMUR01",
|
||||
"user_id": "uuid-user-dari-supabase-auth"
|
||||
}
|
||||
```
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
{
|
||||
"message": "Device berhasil diklaim",
|
||||
"device": {
|
||||
"device_id": "esp32-01",
|
||||
"label": "Kumbung Barat",
|
||||
"location": "Sektor A"
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
#### 2. Ambil Info Device Milik User
|
||||
* **Endpoint**: `GET /api/device/my-device/:userId`
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
{
|
||||
"device_id": "esp32-01",
|
||||
"label": "Kumbung Barat",
|
||||
"location": "Sektor A",
|
||||
"is_online": true,
|
||||
"last_seen": "2026-05-24T10:00:00Z"
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### 🌡️ Threshold Management
|
||||
|
||||
#### 1. Ambil Threshold Aktif
|
||||
* **Endpoint**: `GET /api/threshold/:deviceId`
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
{
|
||||
"device_id": "esp32-01",
|
||||
"temp_max": 30.0,
|
||||
"hum_max": 80.0,
|
||||
"updated_at": "2026-05-24T09:00:00Z"
|
||||
}
|
||||
```
|
||||
|
||||
#### 2. Update Threshold
|
||||
Memperbarui batas threshold di DB dan langsung mengirimkan pembaruan ke ESP32 secara instan via MQTT.
|
||||
* **Endpoint**: `POST /api/threshold/:deviceId`
|
||||
* **Body Request**:
|
||||
```json
|
||||
{
|
||||
"temp_max": 31.5,
|
||||
"hum_max": 85.0
|
||||
}
|
||||
```
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
{
|
||||
"message": "Threshold diupdate",
|
||||
"data": { "device_id": "esp32-01", "temp_max": 31.5, "hum_max": 85.0 }
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### 🗓️ Jadwal & Kontrol Penyiraman
|
||||
|
||||
#### 1. Ambil Semua Jadwal Perangkat
|
||||
* **Endpoint**: `GET /api/schedule/:deviceId`
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
[
|
||||
{
|
||||
"id": "uuid-jadwal-1",
|
||||
"device_id": "esp32-01",
|
||||
"label": "Siram Pagi Hari",
|
||||
"cron": "0 6 * * *",
|
||||
"duration_s": 60,
|
||||
"is_active": true,
|
||||
"created_at": "2026-05-24T08:00:00Z"
|
||||
}
|
||||
]
|
||||
```
|
||||
|
||||
#### 2. Buat Jadwal Penyiraman Baru
|
||||
Menambahkan jadwal ke DB dan mendaftarkannya sebagai antrian *repeatable job* di BullMQ.
|
||||
* **Endpoint**: `POST /api/schedule/:deviceId`
|
||||
* **Body Request**:
|
||||
```json
|
||||
{
|
||||
"label": "Siram Sore",
|
||||
"cron": "0 17 * * *",
|
||||
"duration_s": 45
|
||||
}
|
||||
```
|
||||
* **Respons Sukses (210/201)**:
|
||||
```json
|
||||
{
|
||||
"id": "uuid-jadwal-baru",
|
||||
"device_id": "esp32-01",
|
||||
"label": "Siram Sore",
|
||||
"cron": "0 17 * * *",
|
||||
"duration_s": 45,
|
||||
"is_active": true
|
||||
}
|
||||
```
|
||||
|
||||
#### 3. Hapus Jadwal
|
||||
Menghapus permanen jadwal dari database dan membatalkan repeatable job dari BullMQ.
|
||||
* **Endpoint**: `DELETE /api/schedule/:id`
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
{ "message": "Jadwal dihapus" }
|
||||
```
|
||||
|
||||
#### 4. Toggle Jadwal (Aktif / Nonaktif)
|
||||
Menghentikan eksekusi jadwal sementara waktu di BullMQ tanpa menghapus data jadwal dari database.
|
||||
* **Endpoint**: `PATCH /api/schedule/:id/toggle`
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
{
|
||||
"message": "Jadwal dinonaktifkan",
|
||||
"data": { "id": "uuid-jadwal-1", "is_active": false }
|
||||
}
|
||||
```
|
||||
|
||||
#### 5. Trigger Siram Manual (Sekali Jalan)
|
||||
Memicu penyiraman instan di luar jadwal reguler.
|
||||
* **Endpoint**: `POST /api/schedule/:deviceId/now`
|
||||
* **Body Request** (Opsional):
|
||||
```json
|
||||
{ "duration_s": 30 }
|
||||
```
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
{ "message": "Siram manual 30s dijadwalkan" }
|
||||
```
|
||||
|
||||
#### 6. Hentikan Pompa Secara Paksa
|
||||
Segera mengirimkan sinyal relay `OFF` via MQTT untuk mematikan pompa secara langsung.
|
||||
* **Endpoint**: `POST /api/schedule/:deviceId/stop`
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
{ "message": "Pompa dimatikan" }
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### ⚙️ Mode Operasi Perangkat
|
||||
|
||||
ESP32 mendukung 3 mode operasi utama:
|
||||
* `auto`: Pompa dikendalikan secara otomatis berdasarkan data sensor SHT31.
|
||||
* `manual`: Logika otomatis dimatikan, kontrol pompa sepenuhnya diatur manual lewat API/Aplikasi.
|
||||
* `offline`: Mode darurat jika terputus dari jaringan cloud (logika auto berjalan mandiri di hardware lokal tanpa mempublikasikan data MQTT).
|
||||
|
||||
#### 1. Ambil Mode Kerja Aktif
|
||||
* **Endpoint**: `GET /api/mode/:deviceId`
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
{
|
||||
"device_id": "esp32-01",
|
||||
"current_mode": "auto",
|
||||
"is_online": true,
|
||||
"last_seen": "2026-05-24T10:00:00Z"
|
||||
}
|
||||
```
|
||||
|
||||
#### 2. Ubah Mode Kerja
|
||||
Mengirim perintah ganti mode ke hardware via MQTT dan menyimpan perubahannya di database.
|
||||
* **Endpoint**: `POST /api/mode/:deviceId`
|
||||
* **Body Request**:
|
||||
```json
|
||||
{ "mode": "manual" }
|
||||
```
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
{
|
||||
"message": "Mode berhasil diubah ke manual",
|
||||
"mode": "manual",
|
||||
"changed": true
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### 📊 Riwayat Sensor & Tren Kondisi
|
||||
|
||||
#### 1. Ambil Log Sensor Terbaru
|
||||
Mengambil log real-time sensor terbaru (termasuk filter in-memory deadband).
|
||||
* **Endpoint**: `GET /api/history/:deviceId?limit=100`
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
[
|
||||
{
|
||||
"temperature": 28.5,
|
||||
"humidity": 82.3,
|
||||
"relay_state": false,
|
||||
"mode": "auto",
|
||||
"created_at": "2026-05-24T10:15:00Z"
|
||||
}
|
||||
]
|
||||
```
|
||||
|
||||
#### 2. Ambil Rata-rata Sensor Harian (RPC get_daily_average)
|
||||
Untuk visualisasi grafik jangka panjang di aplikasi Flutter.
|
||||
* **Endpoint**: `GET /api/history/:deviceId/daily?days=7`
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
[
|
||||
{
|
||||
"day": "2026-05-24",
|
||||
"avg_temp": 28.1,
|
||||
"avg_hum": 83.4
|
||||
}
|
||||
]
|
||||
```
|
||||
|
||||
#### 3. Ambil Rata-rata Sensor Per Jam (RPC get_hourly_average)
|
||||
Untuk visualisasi grafik analitis jangka menengah/harian.
|
||||
* **Endpoint**: `GET /api/history/:deviceId/hourly?days=7`
|
||||
* **Respons Sukses (200)**:
|
||||
```json
|
||||
[
|
||||
{
|
||||
"hour": "2026-05-24T10:00:00.000Z",
|
||||
"avg_temp": 28.5,
|
||||
"avg_hum": 82.9
|
||||
}
|
||||
]
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 📡 MQTT Topics
|
||||
|
||||
Sistem komunikasi backend dan ESP32 menggunakan protokol MQTT over TLS (mqtts) di port 8883.
|
||||
|
||||
| Topic | Arah | Payload | Keterangan |
|
||||
|---|---|---|---|
|
||||
| `sensor/sht31` | ESP32 → Backend | `{"device_id":"esp32-01","temp":28.5,"hum":82.3,"mode":"auto","relay":false}` | Laporan status sensor & hardware berkala |
|
||||
| `config/threshold/{deviceId}` | Backend → ESP32 | `{"temp":31.5,"hum":85.0}` | Sinkronisasi perubahan threshold sensor |
|
||||
| `cmd/relay/{deviceId}` | Backend → ESP32 | `"ON"` atau `"OFF"` | Perintah langsung kontrol relay pompa |
|
||||
| `cmd/mode/{deviceId}` | Backend → ESP32 | `"auto"`, `"manual"`, atau `"offline"` | Perintah langsung untuk mengubah mode operasi |
|
||||
|
||||
---
|
||||
|
||||
## 🔁 Antrian Kerja (BullMQ) & Recovery Sistem
|
||||
|
||||
Sistem penjadwalan dikelola secara hybrid menggunakan **BullMQ** dan **Supabase**. Hal ini memecahkan masalah hilangnya repeatable job jika server mengalami restart/deployment ulang.
|
||||
|
||||
### Alur Kerja BullMQ Worker
|
||||
1. Job repeatable terpicu sesuai jadwal *cron* atau instan dari trigger siram manual.
|
||||
2. Worker BullMQ mengambil job dari Redis:
|
||||
- Mengirim perintah relay `ON` ke ESP32 via MQTT.
|
||||
- Mengirimkan push notification "Penyiraman Dimulai" via OneSignal ke pemilik alat.
|
||||
- Menahan eksekusi (*non-blocking sleep*) selama durasi siram `duration_s`.
|
||||
- Mengirim perintah relay `OFF` ke ESP32 via MQTT saat durasi berakhir.
|
||||
|
||||
### Alur Recovery (Schedule Restore)
|
||||
Saat server pertama kali menyala (atau setelah failover berpindah ke aktif):
|
||||
1. Menghapus semua job antrian lama di Redis demi mencegah tumpang tindih.
|
||||
2. Membaca semua data jadwal yang aktif (`is_active = true`) di tabel Supabase.
|
||||
3. Mendaftarkan ulang job ke Redis menggunakan pustaka terbaru **BullMQ v5+** (menggunakan format parameter `cron`).
|
||||
|
||||
---
|
||||
|
||||
## 📁 Struktur Folder
|
||||
|
||||
```
|
||||
backend-jamur/
|
||||
├── src/
|
||||
│ ├── index.js # Entry point utama, Express & rate limit setup
|
||||
│ ├── jobs/
|
||||
│ │ └── offlineDetector.js# Background job pemeriksa status online perangkat
|
||||
│ ├── mqtt/
|
||||
│ │ └── mqttClient.js # Koneksi HiveMQ, subscriber sensor, & publisher command
|
||||
│ ├── queues/
|
||||
│ │ ├── irrigationQueue.js# Inisialisasi antrian BullMQ & koneksi Redis Upstash
|
||||
│ │ ├── irrigationWorker.js# Worker pengeksekusi siklus pompa ON -> DELAY -> OFF
|
||||
│ │ └── scheduleRestore.js# Pemulihan otomatis jadwal aktif dari database saat startup
|
||||
│ ├── routes/
|
||||
│ │ ├── device.js # API endpoint klaim & info perangkat
|
||||
│ │ ├── history.js # API endpoint riwayat & agregasi sensor (daily/hourly)
|
||||
│ │ ├── mode.js # API endpoint kontrol mode kerja ESP32
|
||||
│ │ ├── schedule.js # API endpoint CRUD jadwal & trigger pompa manual
|
||||
│ │ └── threshold.js # API endpoint pembacaan & modifikasi batas sensor
|
||||
│ ├── supabase/
|
||||
│ │ └── client.js # Setup Supabase JS client (Service Role authorization)
|
||||
│ └── utils/
|
||||
│ └── notification.js # Pengirim OneSignal Push Notification + Cooldown redis
|
||||
├── .env # File konfigurasi privat (lokal)
|
||||
├── .gitignore # Daftar file terabaikan dari Git
|
||||
├── package.json # Daftar pustaka dependencies & scripts npm
|
||||
└── package-lock.json # Lock file dependencies
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 📦 Dependensi Utama
|
||||
|
||||
Detail library penting yang digunakan pada proyek ini:
|
||||
|
||||
```json
|
||||
"dependencies": {
|
||||
"@supabase/supabase-js": "^2.101.1",
|
||||
"bullmq": "^5.73.0",
|
||||
"dotenv": "^17.4.1",
|
||||
"express": "^5.2.1",
|
||||
"express-rate-limit": "^8.3.2",
|
||||
"ioredis": "^5.10.1",
|
||||
"mqtt": "^5.15.1"
|
||||
}
|
||||
```
|
||||
File diff suppressed because it is too large
Load Diff
|
|
@ -0,0 +1,25 @@
|
|||
{
|
||||
"name": "broker",
|
||||
"version": "1.0.0",
|
||||
"main": "index.js",
|
||||
"scripts": {
|
||||
"test": "echo \"Error: no test specified\" && exit 1",
|
||||
"start": "node src/index.js",
|
||||
"dev": "nodemon src/index.js"
|
||||
},
|
||||
"author": "",
|
||||
"license": "ISC",
|
||||
"description": "",
|
||||
"dependencies": {
|
||||
"@supabase/supabase-js": "^2.101.1",
|
||||
"bullmq": "^5.73.0",
|
||||
"dotenv": "^17.4.1",
|
||||
"express": "^5.2.1",
|
||||
"express-rate-limit": "^8.3.2",
|
||||
"ioredis": "^5.10.1",
|
||||
"mqtt": "^5.15.1"
|
||||
},
|
||||
"devDependencies": {
|
||||
"nodemon": "^3.1.14"
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,98 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* ENTRY POINT — Titik Masuk Aplikasi Backend
|
||||
* =============================================================================
|
||||
* File PERTAMA yang berjalan saat server dinyalakan.
|
||||
* Tugasnya: setup semua komponen (Express, MQTT, Worker, Rate Limiter)
|
||||
* dan menghubungkan semuanya menjadi satu aplikasi yang siap melayani request.
|
||||
*
|
||||
* PENTING — Urutan startup sangat penting:
|
||||
* 1. connect() -> Koneksi MQTT dijalankan PERTAMA karena Worker butuh MQTT.
|
||||
* 2. require Worker -> Worker diaktifkan SETELAH MQTT siap, agar tidak error
|
||||
* ketika Worker mencoba publishRelay di job pertamanya.
|
||||
* 3. restoreSchedules() -> Baca jadwal aktif dari DB dan daftarkan ulang ke BullMQ.
|
||||
* Ini menjamin jadwal tidak hilang setelah server restart/redeploy.
|
||||
* 4. Baru setelah itu Express & semua route dijalankan.
|
||||
*
|
||||
* Kalau urutannya salah (misal Worker dijalankan sebelum MQTT connect),
|
||||
* penyiraman pertama bisa gagal karena client MQTT belum siap.
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
require('dotenv').config() // Load semua variabel dari file .env ke process.env
|
||||
const express = require('express')
|
||||
const rateLimit = require('express-rate-limit')
|
||||
const { initFailover } = require('./utils/failoverManager')
|
||||
|
||||
// Inisialisasi failover dinamis (menggantikan startup statis Langkah 1-4)
|
||||
initFailover()
|
||||
|
||||
// LANGKAH 5: Inisialisasi aplikasi Express
|
||||
const app = express()
|
||||
|
||||
// Percayai layer proxy di depan aplikasi secara dinamis.
|
||||
// Jika di VPS (Backup) dibatasi 1 layer proxy (Nginx).
|
||||
// Jika di Railway (Primary) dibatasi 2 layer proxy (Nginx di VPS -> Edge Proxy Railway -> Node.js).
|
||||
const trustProxyLayers = process.env.IS_BACKUP_SERVER === 'true' ? 1 : 2
|
||||
app.set('trust proxy', trustProxyLayers)
|
||||
|
||||
app.use(express.json()) // Supaya server bisa membaca body request dalam format JSON
|
||||
|
||||
// Middleware logger sederhana untuk membantu memantau request di PM2 logs (terutama saat failover)
|
||||
app.use((req, res, next) => {
|
||||
console.log(`[HTTP] ${req.method} ${req.url}`)
|
||||
next()
|
||||
})
|
||||
|
||||
// =============================================================================
|
||||
// RATE LIMITER — Pelindung dari request berlebihan
|
||||
// =============================================================================
|
||||
// Rate limiter global: berlaku untuk semua endpoint /api
|
||||
// Maksimal 100 request per menit dari satu IP address yang sama
|
||||
const globalLimiter = rateLimit({
|
||||
windowMs: 60 * 1000, // Window waktu: 1 menit
|
||||
max: 100, // Maks 100 request per window
|
||||
standardHeaders: true,
|
||||
legacyHeaders: false,
|
||||
message: { error: 'Terlalu banyak request, coba lagi dalam 1 menit.' },
|
||||
})
|
||||
|
||||
// Rate limiter ketat khusus untuk endpoint siram manual
|
||||
// Dibatasi hanya 5x per menit untuk mencegah penyiraman berlebihan yang
|
||||
// bisa merusak tanaman atau menguras air terlalu cepat
|
||||
const manualIrrigationLimiter = rateLimit({
|
||||
windowMs: 60 * 1000, // Window waktu: 1 menit
|
||||
max: 5, // Maks 5 kali siram manual per menit
|
||||
standardHeaders: true,
|
||||
legacyHeaders: false,
|
||||
message: { error: 'Terlalu sering menyiram, tunggu sebentar.' },
|
||||
})
|
||||
// =============================================================================
|
||||
|
||||
// Terapkan rate limiter global ke semua route /api
|
||||
app.use('/api', globalLimiter)
|
||||
|
||||
// =============================================================================
|
||||
// ROUTE REGISTRATION — Daftarkan semua endpoint API
|
||||
// =============================================================================
|
||||
app.use('/api/threshold', require('./routes/threshold')) // Atur batas suhu & kelembapan
|
||||
app.use('/api/history', require('./routes/history')) // Lihat riwayat data sensor
|
||||
app.use('/api/schedule/:deviceId/now', manualIrrigationLimiter) // Extra ketat untuk siram manual (HARUS sebelum route schedule!)
|
||||
app.use('/api/schedule', require('./routes/schedule')) // Kelola jadwal penyiraman
|
||||
app.use('/api/device', require('./routes/device')) // Klaim & kelola perangkat
|
||||
app.use('/api/mode', require('./routes/mode')) // Kontrol mode operasi ESP32 (auto/manual/offline)
|
||||
app.use('/api/latency', require('./routes/latency')) // Tes latensi dari backend asli ke MQTT broker
|
||||
// =============================================================================
|
||||
|
||||
// Endpoint health check untuk platform deployment seperti Fly.io / Railway
|
||||
app.get('/', (req, res) => {
|
||||
res.status(200).send('OK');
|
||||
});
|
||||
|
||||
// Railway secara otomatis menyediakan variabel PORT via environment variable.
|
||||
// Jangan hardcode port! Gunakan process.env.PORT agar server bisa berjalan di Railway.
|
||||
// Fallback ke 3000 hanya untuk development lokal.
|
||||
const PORT = process.env.PORT || 3000
|
||||
app.listen(PORT, () => {
|
||||
console.log(`[Server] Running on port ${PORT}`)
|
||||
})
|
||||
|
|
@ -0,0 +1,63 @@
|
|||
const supabase = require('../supabase/client');
|
||||
const { sendNotification } = require('../utils/notification');
|
||||
|
||||
// Jalankan pengecekan setiap 5 menit
|
||||
const CHECK_INTERVAL_MS = 5 * 60 * 1000;
|
||||
|
||||
let intervalId = null;
|
||||
|
||||
function startOfflineDetector() {
|
||||
if (intervalId) return; // Menghindari multiple intervals
|
||||
intervalId = setInterval(async () => {
|
||||
try {
|
||||
// console.log('[OfflineDetector] Mengecek perangkat offline...');
|
||||
|
||||
// Waktu 5 menit yang lalu (di Supabase pakai UTC)
|
||||
const fiveMinsAgo = new Date(Date.now() - CHECK_INTERVAL_MS).toISOString();
|
||||
|
||||
// Cari perangkat yang is_online = true TAPI last_seen < 5 menit yang lalu
|
||||
const { data: devices, error } = await supabase
|
||||
.from('devices')
|
||||
.select('device_id, claimed_by, last_seen')
|
||||
.eq('is_online', true)
|
||||
.lt('last_seen', fiveMinsAgo);
|
||||
|
||||
if (error) throw error;
|
||||
|
||||
for (const device of devices) {
|
||||
console.log(`[OfflineDetector] Perangkat ${device.device_id} terdeteksi OFFLINE!`);
|
||||
|
||||
// 1. Update status di DB jadi offline
|
||||
await supabase
|
||||
.from('devices')
|
||||
.update({ is_online: false })
|
||||
.eq('device_id', device.device_id);
|
||||
|
||||
// 2. Kirim Notifikasi OneSignal
|
||||
// Cooldown 1 jam (3600 detik) di-handle otomatis oleh sendNotification
|
||||
if (device.claimed_by) {
|
||||
sendNotification(
|
||||
device.claimed_by,
|
||||
'Perangkat Offline! 🚨',
|
||||
`Perangkat IoT Kumbung (${device.device_id}) terputus dari jaringan atau mati listrik.`,
|
||||
3600
|
||||
);
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
console.error('[OfflineDetector] Error:', error.message);
|
||||
}
|
||||
}, CHECK_INTERVAL_MS);
|
||||
|
||||
console.log('[OfflineDetector] Aktif. Pengecekan perangkat offline berjalan setiap 5 menit.');
|
||||
}
|
||||
|
||||
function stopOfflineDetector() {
|
||||
if (intervalId) {
|
||||
clearInterval(intervalId);
|
||||
intervalId = null;
|
||||
console.log('[OfflineDetector] Dinonaktifkan.');
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = { startOfflineDetector, stopOfflineDetector };
|
||||
|
|
@ -0,0 +1,442 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* MQTT CLIENT — Jembatan Komunikasi ke ESP32
|
||||
* =============================================================================
|
||||
* File ini mengurus semua komunikasi antara server Node.js dan perangkat ESP32
|
||||
* menggunakan protokol MQTT melalui HiveMQ Cloud.
|
||||
*
|
||||
* MQTT itu apa?
|
||||
* MQTT adalah protokol pesan ringan yang populer di dunia IoT. Cara kerjanya
|
||||
* seperti sistem "publish-subscribe": ada yang kirim pesan (publish) ke sebuah
|
||||
* "topik", dan ada yang langganan (subscribe) topik itu untuk menerimanya.
|
||||
*
|
||||
* Di sistem ini:
|
||||
* - ESP32 → PUBLISH data sensor ke topik 'sensor/sht31'
|
||||
* - Server → SUBSCRIBE ke 'sensor/sht31' untuk menerima data tersebut
|
||||
* - Server → PUBLISH perintah ke 'cmd/relay/{deviceId}' untuk nyala/matiin pompa
|
||||
* - Server → PUBLISH config ke 'config/threshold' saat threshold diubah
|
||||
*
|
||||
* Fitur penting di file ini: IN-MEMORY CACHE untuk threshold.
|
||||
* Tanpa cache: setiap data sensor masuk (bisa 1x/detik!) → query ke Supabase.
|
||||
* Dengan cache: data threshold disimpan di memori, query DB hanya tiap 30 detik.
|
||||
* Ini menghemat kuota database secara signifikan.
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
const mqtt = require('mqtt')
|
||||
const supabase = require('../supabase/client')
|
||||
const { sendNotification } = require('../utils/notification')
|
||||
|
||||
// Daftar topik MQTT yang digunakan dalam sistem ini
|
||||
const TOPIC_SENSOR = 'sensor/sht31' // Topik untuk menerima data dari ESP32
|
||||
const TOPIC_THRESHOLD_BASE = 'config/threshold' // Topik dasar untuk kirim setting threshold ke ESP32
|
||||
const TOPIC_RELAY = 'cmd/relay' // Topik dasar untuk kontrol relay
|
||||
const TOPIC_MODE = 'cmd/mode' // Topik dasar untuk mengirim perintah ganti mode ke ESP32
|
||||
|
||||
// =============================================================================
|
||||
// IN-MEMORY CACHE — Penyimpanan Sementara di Memori Server
|
||||
// =============================================================================
|
||||
// Map ini menyimpan data threshold tiap device agar tidak perlu query DB terus-menerus.
|
||||
// Format isi Map: { deviceId → { temp_min, temp_max, hum_max, cachedAt } }
|
||||
const thresholdCache = new Map()
|
||||
const CACHE_TTL_MS = 30 * 1000 // 30 detik dalam milidetik
|
||||
|
||||
/**
|
||||
* Ambil data threshold dari cache. Kalau cache kosong atau sudah kadaluarsa,
|
||||
* baru query ke database Supabase, lalu simpan hasilnya ke cache.
|
||||
*
|
||||
* @param {string} deviceId - ID perangkat ESP32
|
||||
* @returns {Object|null} Objek { temp_min, temp_max, hum_max } atau null jika gagal
|
||||
*/
|
||||
async function getCachedThreshold(deviceId) {
|
||||
const cached = thresholdCache.get(deviceId)
|
||||
|
||||
// Cek apakah cache masih valid (belum kadaluarsa)
|
||||
if (cached && (Date.now() - cached.cachedAt) < CACHE_TTL_MS) {
|
||||
return cached // cache hit — tidak perlu query DB
|
||||
}
|
||||
|
||||
// Cache miss atau sudah expired — ambil data segar dari Supabase
|
||||
const { data, error } = await supabase
|
||||
.from('thresholds')
|
||||
.select('temp_min, temp_max, hum_max')
|
||||
.eq('device_id', deviceId)
|
||||
.limit(1)
|
||||
.maybeSingle()
|
||||
|
||||
if (error) {
|
||||
console.error('[Cache] Gagal baca threshold:', error.message)
|
||||
return null
|
||||
}
|
||||
|
||||
if (!data) {
|
||||
console.warn(`[Cache] Threshold belum disetel untuk ${deviceId}`)
|
||||
return null
|
||||
}
|
||||
|
||||
// Fallback temp_min ke 20.0 jika di db nilainya null
|
||||
const temp_min = (data.temp_min === undefined || data.temp_min === null) ? 20.0 : data.temp_min;
|
||||
|
||||
// Simpan hasil query ke cache beserta timestamp saat ini
|
||||
thresholdCache.set(deviceId, {
|
||||
temp_min,
|
||||
temp_max: data.temp_max,
|
||||
hum_max: data.hum_max,
|
||||
cachedAt: Date.now()
|
||||
})
|
||||
console.log(`[Cache] Threshold ${deviceId} diperbarui dari DB`)
|
||||
return { temp_min, temp_max: data.temp_max, hum_max: data.hum_max }
|
||||
}
|
||||
// =============================================================================
|
||||
|
||||
// =============================================================================
|
||||
// IN-MEMORY CACHE UNTUK SMART FILTER (DEADBAND) & STATUS
|
||||
// =============================================================================
|
||||
// Map ini mencatat status terakhir agar tidak spam ke Supabase
|
||||
const sensorLogCache = new Map() // { temp, hum, relay_state, mode, lastSavedAt }
|
||||
const lastSeenCache = new Map() // timestamp terakhir kali device laporan online
|
||||
const pendingPings = new Map() // Menyimpan callback untuk pencocokan pingId latency test
|
||||
|
||||
const TEMP_DELTA = 0.5; // Perubahan suhu minimal untuk disimpan
|
||||
const HUM_DELTA = 1.0; // Perubahan kelembapan minimal untuk disimpan
|
||||
const HEARTBEAT_INTERVAL_MS = 5 * 60 * 1000; // 5 menit (jaga grafik tidak putus)
|
||||
const LAST_SEEN_INTERVAL_MS = 1 * 60 * 1000; // 1 menit (throttle status online)
|
||||
// =============================================================================
|
||||
|
||||
// Variabel untuk menyimpan instance koneksi MQTT (digunakan di seluruh file ini)
|
||||
let client
|
||||
|
||||
/**
|
||||
* Memulai koneksi ke MQTT broker (HiveMQ Cloud).
|
||||
* Fungsi ini dipanggil PERTAMA KALI saat server start di index.js,
|
||||
* sebelum Worker dijalankan, agar MQTT siap saat Worker butuh publishRelay.
|
||||
*/
|
||||
function connect() {
|
||||
if (client) {
|
||||
console.warn('[MQTT] Client sudah terhubung. Mengabaikan pemanggilan connect() baru.')
|
||||
return
|
||||
}
|
||||
// Buat koneksi ke HiveMQ menggunakan kredensial dari .env
|
||||
client = mqtt.connect(process.env.MQTT_BROKER_URL, {
|
||||
port: parseInt(process.env.MQTT_PORT) || 8883,
|
||||
username: process.env.MQTT_USERNAME,
|
||||
password: process.env.MQTT_PASSWORD,
|
||||
protocol: 'mqtts', // wajib TLS untuk HiveMQ Cloud (koneksi terenkripsi)
|
||||
clientId: `nodejs-backend-${Date.now()}`, // ID unik agar tidak bentrok
|
||||
clean: true,
|
||||
reconnectPeriod: 5000, // Coba reconnect tiap 5 detik jika koneksi putus
|
||||
})
|
||||
|
||||
// Event: Berhasil terhubung ke broker
|
||||
client.on('connect', () => {
|
||||
console.log('[MQTT] Terhubung ke HiveMQ Cloud')
|
||||
// Mulai "dengarkan" data sensor dari ESP32
|
||||
client.subscribe(TOPIC_SENSOR, { qos: 1 })
|
||||
// Mulai "dengarkan" status keaktifan perangkat (LWT)
|
||||
client.subscribe('status/+', { qos: 1 })
|
||||
// Mulai "dengarkan" pesan ping latensi
|
||||
client.subscribe('latency/test/ping', { qos: 1 })
|
||||
})
|
||||
|
||||
/**
|
||||
* Event: Ada pesan masuk dari topik yang di-subscribe.
|
||||
* Ini adalah inti dari pemrosesan data sensor.
|
||||
* Alur: Terima data → Simpan ke DB → Cek threshold → Kontrol relay
|
||||
*/
|
||||
client.on('message', async (topic, payload) => {
|
||||
try {
|
||||
// Proses ping latensi loopback backend
|
||||
if (topic === 'latency/test/ping') {
|
||||
try {
|
||||
const data = JSON.parse(payload.toString())
|
||||
const { pingId } = data
|
||||
const callback = pendingPings.get(pingId)
|
||||
if (callback) {
|
||||
callback(Date.now())
|
||||
pendingPings.delete(pingId)
|
||||
}
|
||||
} catch (e) {}
|
||||
return
|
||||
}
|
||||
|
||||
// Proses pesan status LWT (status/+)
|
||||
if (topic.startsWith('status/')) {
|
||||
const device_id = topic.split('/')[1]
|
||||
const status = payload.toString().trim() // "online" atau "offline"
|
||||
console.log(`[MQTT] [LWT] Perangkat ${device_id} berstatus: ${status}`)
|
||||
|
||||
const isOnline = status === 'online'
|
||||
|
||||
// Update status di Supabase
|
||||
const { data, error } = await supabase.from('devices')
|
||||
.update({
|
||||
is_online: isOnline,
|
||||
last_seen: new Date().toISOString()
|
||||
})
|
||||
.eq('device_id', device_id)
|
||||
.select('claimed_by')
|
||||
.maybeSingle()
|
||||
|
||||
if (error) {
|
||||
console.error(`[MQTT] [LWT] Gagal update status ${device_id} ke DB:`, error.message)
|
||||
} else {
|
||||
console.log(`[MQTT] [LWT] Berhasil update status ${device_id} ke DB -> is_online: ${isOnline}`)
|
||||
if (!isOnline && data && data.claimed_by) {
|
||||
sendNotification(
|
||||
data.claimed_by,
|
||||
'Perangkat Offline! 🚨',
|
||||
`Perangkat IoT Kumbung (${device_id}) terputus dari jaringan atau mati listrik (LWT).`,
|
||||
3600 // Cooldown 1 jam
|
||||
);
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
// Hanya proses pesan dari topik sensor, abaikan topik lain
|
||||
if (topic !== TOPIC_SENSOR) return
|
||||
|
||||
// Parsing payload dari format JSON ke objek JavaScript
|
||||
let data
|
||||
try {
|
||||
data = JSON.parse(payload.toString())
|
||||
} catch {
|
||||
return console.error('[MQTT] Payload tidak valid JSON')
|
||||
}
|
||||
|
||||
// Destructure data sensor. Jika ESP32 tidak kirim device_id, pakai default 'esp32-01'
|
||||
// relay dikirim ESP32 sebagai boolean (true = ON, false = OFF) pada key 'relay' atau 'relay_state'
|
||||
// mode dikirim ESP32 sebagai string (contoh: "auto", "auto-on", "cooldown", "schedule-on")
|
||||
const { temp, hum, mode, device_id = 'esp32-01' } = data
|
||||
const relay_state = data.relay_state ?? data.relay ?? false;
|
||||
console.log(`[MQTT] [${device_id}] Sensor: ${temp}°C | ${hum}% | Relay: ${relay_state}`)
|
||||
|
||||
const now = Date.now();
|
||||
|
||||
// Langkah 1A: Update status online (last_seen) dengan throttle 1 menit
|
||||
const lastSeen = lastSeenCache.get(device_id) || 0;
|
||||
if (now - lastSeen > LAST_SEEN_INTERVAL_MS) {
|
||||
// Kita tidak perlu await di sini agar tidak memblokir proses peringatan (Push Notif)
|
||||
supabase.from('devices')
|
||||
.update({ is_online: true, last_seen: new Date().toISOString() })
|
||||
.eq('device_id', device_id)
|
||||
.then(({ error }) => {
|
||||
if (error) console.error('[MQTT] Gagal update last_seen:', error.message);
|
||||
else lastSeenCache.set(device_id, now);
|
||||
});
|
||||
}
|
||||
|
||||
// Langkah 1B: Smart Filter untuk menyimpan data ke tabel sensor_logs
|
||||
let shouldSaveData = false;
|
||||
const lastData = sensorLogCache.get(device_id);
|
||||
|
||||
// Pastikan mode tidak pernah null agar Flutter app tidak crash (terpental)
|
||||
const safeMode = mode ?? (lastData?.mode ?? 'auto');
|
||||
|
||||
if (!lastData) {
|
||||
shouldSaveData = true; // Simpan jika belum ada riwayat di memori server
|
||||
} else {
|
||||
const tempDiff = Math.abs(temp - lastData.temp);
|
||||
const humDiff = Math.abs(hum - lastData.hum);
|
||||
const timeDiff = now - lastData.lastSavedAt;
|
||||
const relayChanged = relay_state !== lastData.relay_state;
|
||||
const modeChanged = safeMode !== lastData.mode;
|
||||
|
||||
// Logika Smart Filter (Deadband)
|
||||
if (tempDiff > TEMP_DELTA || humDiff > HUM_DELTA || relayChanged || modeChanged || timeDiff > HEARTBEAT_INTERVAL_MS) {
|
||||
shouldSaveData = true;
|
||||
}
|
||||
}
|
||||
|
||||
if (shouldSaveData) {
|
||||
const { error: insertErr } = await supabase.from('sensor_logs').insert({
|
||||
device_id,
|
||||
temperature: temp,
|
||||
humidity: hum,
|
||||
relay_state: relay_state ?? false,
|
||||
mode: safeMode, // Simpan snapshot mode ESP32 saat data dikirim
|
||||
})
|
||||
if (insertErr) {
|
||||
console.error('[MQTT] Gagal simpan sensor log:', insertErr.message)
|
||||
} else {
|
||||
// Catat ke memori setelah berhasil insert
|
||||
sensorLogCache.set(device_id, {
|
||||
temp,
|
||||
hum,
|
||||
relay_state: relay_state ?? false,
|
||||
mode: safeMode,
|
||||
lastSavedAt: now
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
// Langkah 2: Baca threshold dari cache (hanya untuk memastikan cache tersimpan, logika auto diurus ESP32)
|
||||
const threshold = await getCachedThreshold(device_id)
|
||||
|
||||
// Langkah 3: Pengecekan Suhu & Kelembapan untuk Push Notification
|
||||
// Anti-spam (cooldown 30 menit) kini diurus otomatis oleh sendNotification.
|
||||
if (threshold) {
|
||||
const alerts = []
|
||||
let isCold = false
|
||||
let isHot = false
|
||||
let isDry = false
|
||||
|
||||
if (hum < threshold.hum_max) {
|
||||
alerts.push(`Kelembapan saat ini ${hum}% (Batas minimal: ${threshold.hum_max}%)`)
|
||||
isDry = true
|
||||
}
|
||||
if (temp > threshold.temp_max) {
|
||||
alerts.push(`Suhu saat ini ${temp}°C (Batas maksimal: ${threshold.temp_max}°C)`)
|
||||
isHot = true
|
||||
} else if (temp < threshold.temp_min) {
|
||||
alerts.push(`Suhu saat ini ${temp}°C (Batas minimal: ${threshold.temp_min}°C)`)
|
||||
isCold = true
|
||||
}
|
||||
|
||||
if (alerts.length > 0) {
|
||||
let alertMsg = ''
|
||||
if (alerts.length >= 2) {
|
||||
alertMsg = `Peringatan Kumbung! ${alerts.join(' dan ')}`
|
||||
} else if (isCold) {
|
||||
alertMsg = `Peringatan Dingin! ${alerts[0]}`
|
||||
} else if (isDry) {
|
||||
alertMsg = `Peringatan Kering! ${alerts[0]}`
|
||||
} else if (isHot) {
|
||||
alertMsg = `Peringatan Panas! ${alerts[0]}`
|
||||
}
|
||||
|
||||
// Cari user yang memiliki alat ini
|
||||
const { data: device } = await supabase
|
||||
.from('devices')
|
||||
.select('claimed_by')
|
||||
.eq('device_id', device_id)
|
||||
.single()
|
||||
|
||||
if (device && device.claimed_by) {
|
||||
// Cooldown 30 menit (1800 detik) di-handle oleh sendNotification
|
||||
sendNotification(device.claimed_by, 'Peringatan Sensor Kumbung ⚠️', alertMsg, 1800)
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (globalErr) {
|
||||
console.error('[MQTT] Kesalahan tidak terduga saat memproses pesan:', globalErr.message)
|
||||
}
|
||||
})
|
||||
|
||||
// Event: Terjadi error pada koneksi MQTT
|
||||
client.on('error', err => console.error('[MQTT] Error:', err.message))
|
||||
|
||||
// Event: Koneksi terputus — otomatis akan mencoba reconnect sesuai reconnectPeriod
|
||||
client.on('offline', () => console.warn('[MQTT] Koneksi terputus, mencoba reconnect...'))
|
||||
}
|
||||
|
||||
/**
|
||||
* Memutuskan koneksi MQTT secara aman (digunakan saat failover kembali ke standby).
|
||||
*/
|
||||
function disconnect() {
|
||||
if (client) {
|
||||
client.end(true, () => {
|
||||
console.log('[MQTT] Terputus secara aman (dinonaktifkan oleh failover)')
|
||||
})
|
||||
client = null
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Kirim nilai threshold terbaru ke ESP32 via MQTT.
|
||||
* Dipanggil otomatis saat user update threshold melalui API.
|
||||
* Sekaligus memperbarui cache agar nilai baru langsung efektif.
|
||||
*
|
||||
* @param {string} deviceId - ID perangkat target
|
||||
* @param {number} tempMin - Batas minimum suhu (°C)
|
||||
* @param {number} tempMax - Batas maksimum suhu (°C)
|
||||
* @param {number} humMax - Batas maksimum kelembapan (%)
|
||||
*/
|
||||
function publishThreshold(deviceId, tempMin, tempMax, humMax) {
|
||||
const payload = JSON.stringify({
|
||||
temp_min: tempMin,
|
||||
temp_max: tempMax,
|
||||
temp: tempMax, // backward compatibility
|
||||
hum: humMax
|
||||
})
|
||||
// retain: true → ESP32 yang baru connect akan langsung dapat nilai threshold terkini
|
||||
if (client) {
|
||||
client.publish(`${TOPIC_THRESHOLD_BASE}/${deviceId}`, payload, { qos: 1, retain: true })
|
||||
} else {
|
||||
console.warn(`[MQTT] Terputus. Gagal publish threshold untuk ${deviceId} (Server Standby)`)
|
||||
}
|
||||
|
||||
// Perbarui cache langsung agar nilai baru efektif tanpa harus tunggu 30 detik
|
||||
thresholdCache.set(deviceId, { temp_min: tempMin, temp_max: tempMax, hum_max: humMax, cachedAt: Date.now() })
|
||||
console.log(`[Cache] Threshold ${deviceId} diperbarui dari API`)
|
||||
}
|
||||
|
||||
/**
|
||||
* Kirim perintah nyala atau mati ke relay ESP32.
|
||||
* Topik yang digunakan bersifat dinamis per device: cmd/relay/{deviceId}
|
||||
* Sehingga jika ada banyak device, perintah tidak tercampur.
|
||||
*
|
||||
* @param {string} deviceId - ID perangkat target
|
||||
* @param {'ON'|'OFF'} state - Perintah yang dikirim
|
||||
*/
|
||||
function publishRelay(deviceId, state) {
|
||||
if (client) {
|
||||
client.publish(`${TOPIC_RELAY}/${deviceId}`, state, { qos: 1 })
|
||||
console.log(`[MQTT] Relay ${deviceId} → ${state}`)
|
||||
} else {
|
||||
console.warn(`[MQTT] Terputus. Gagal kirim perintah relay ${state} ke ${deviceId} (Server Standby)`)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Kirim perintah ganti mode operasi ke ESP32.
|
||||
* ESP32 akan berpindah ke mode target dan menjalankan logika yang sesuai.
|
||||
*
|
||||
* @param {string} deviceId - ID perangkat target
|
||||
* @param {'auto'|'manual'|'offline'} mode - Mode yang ingin diaktifkan
|
||||
*/
|
||||
function publishMode(deviceId, mode) {
|
||||
if (client) {
|
||||
client.publish(`${TOPIC_MODE}/${deviceId}`, mode, { qos: 1 })
|
||||
console.log(`[MQTT] Mode ${deviceId} → ${mode}`)
|
||||
} else {
|
||||
console.warn(`[MQTT] Terputus. Gagal kirim perintah mode ${mode} ke ${deviceId} (Server Standby)`)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Mengukur latensi RTT (Round Trip Time) antara server backend ini dengan MQTT Broker HiveMQ Cloud.
|
||||
* @returns {Promise<number>} Latensi RTT dalam milidetik
|
||||
*/
|
||||
function measureLatency() {
|
||||
return new Promise((resolve, reject) => {
|
||||
if (!client || !client.connected) {
|
||||
return reject(new Error('Klien MQTT tidak terhubung ke broker'))
|
||||
}
|
||||
|
||||
const pingId = Math.random().toString(36).substring(7)
|
||||
const startTime = Date.now()
|
||||
|
||||
// Batasi tunggu maksimal 5 detik
|
||||
const timeout = setTimeout(() => {
|
||||
pendingPings.delete(pingId)
|
||||
reject(new Error('Pengujian latensi MQTT timeout (5 detik)'))
|
||||
}, 5000)
|
||||
|
||||
pendingPings.set(pingId, (endTime) => {
|
||||
clearTimeout(timeout)
|
||||
resolve(endTime - startTime)
|
||||
})
|
||||
|
||||
const payload = JSON.stringify({ pingId })
|
||||
client.publish('latency/test/ping', payload, { qos: 1 }, (err) => {
|
||||
if (err) {
|
||||
clearTimeout(timeout)
|
||||
pendingPings.delete(pingId)
|
||||
reject(err)
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
module.exports = { connect, disconnect, publishThreshold, publishRelay, publishMode, measureLatency }
|
||||
|
|
@ -0,0 +1,34 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* IRRIGATION QUEUE — Antrian Penyiraman
|
||||
* =============================================================================
|
||||
* File ini mendefinisikan "antrian" penyiraman menggunakan BullMQ.
|
||||
*
|
||||
* Analogi sederhana:
|
||||
* Bayangkan antrian ini seperti "daftar tugas" yang disimpan di Redis.
|
||||
* Saat kamu minta "siram sekarang" atau "jadwalkan siram jam 6 pagi",
|
||||
* permintaan itu tidak langsung dieksekusi — tapi dulu didaftarkan ke antrian ini.
|
||||
* Nanti si Worker (irrigationWorker.js) yang akan mengambil dan mengeksekusinya.
|
||||
*
|
||||
* Kenapa pakai antrian dan tidak langsung eksekusi saja?
|
||||
* - Lebih aman: kalau server mati tiba-tiba, job tidak hilang karena tersimpan di Redis.
|
||||
* - Lebih teratur: tidak ada dua job siram yang jalan bersamaan secara tidak terkontrol.
|
||||
* - Mendukung jadwal (cron): BullMQ bisa menjalankan tugas berulang sesuai pola waktu.
|
||||
*
|
||||
* File ini di-import oleh:
|
||||
* - routes/schedule.js — untuk mendaftarkan jadwal baru atau siram manual
|
||||
* - irrigationWorker.js — untuk tahu ke queue mana ia harus "mendengarkan"
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
const { Queue } = require('bullmq')
|
||||
|
||||
// Koneksi ke Redis menggunakan URL dari .env
|
||||
// Redis berfungsi sebagai "otak" BullMQ: tempat menyimpan semua data antrian
|
||||
const connection = { url: process.env.REDIS_URL }
|
||||
|
||||
// Buat queue dengan nama 'irrigation'
|
||||
// Nama ini harus sama persis dengan yang dipakai di Worker!
|
||||
const irrigationQueue = new Queue('irrigation', { connection })
|
||||
|
||||
module.exports = { irrigationQueue, connection }
|
||||
|
|
@ -0,0 +1,103 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* IRRIGATION WORKER — Eksekutor Penyiraman
|
||||
* =============================================================================
|
||||
* File ini adalah "pekerja" yang secara aktif memantau antrian irrigationQueue
|
||||
* dan mengeksekusi setiap job penyiraman yang masuk.
|
||||
*
|
||||
* Cara kerjanya:
|
||||
* 1. Worker terus "mendengarkan" antrian 'irrigation' di Redis.
|
||||
* 2. Saat ada job masuk (baik dari jadwal cron maupun trigger manual),
|
||||
* Worker langsung mengambil dan menjalankannya.
|
||||
* 3. Pertama, Worker kirim perintah relay ON ke ESP32 via MQTT -> pompa menyala.
|
||||
* 4. Worker menunggu selama `durationSeconds` detik (waktu penyiraman).
|
||||
* 5. Setelah waktu habis, Worker kirim perintah relay OFF -> pompa mati.
|
||||
*
|
||||
* Penting: Worker harus diinisialisasi SETELAH MQTT sudah terkoneksi,
|
||||
* karena worker butuh fungsi publishRelay dari mqttClient.
|
||||
* Urutan startup diatur di src/index.js.
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
const { Worker } = require('bullmq')
|
||||
const { connection } = require('./irrigationQueue')
|
||||
const { publishRelay } = require('../mqtt/mqttClient')
|
||||
const { sendNotification } = require('../utils/notification')
|
||||
const supabase = require('../supabase/client')
|
||||
|
||||
/**
|
||||
* Definisi Worker untuk queue 'irrigation'.
|
||||
* Setiap kali ada job masuk, fungsi async di bawah ini akan dijalankan.
|
||||
*
|
||||
* @param {Object} job - Objek job dari BullMQ
|
||||
* @param {string} job.data.deviceId - ID perangkat ESP32 yang akan disiram
|
||||
* @param {number} job.data.durationSeconds - Lama penyiraman dalam detik
|
||||
*/
|
||||
const worker = new Worker('irrigation', async (job) => {
|
||||
const { deviceId, durationSeconds } = job.data
|
||||
const nowUTC = new Date().toISOString()
|
||||
const nowWIB = new Date(Date.now() + 7 * 60 * 60 * 1000).toISOString().replace('T', ' ').substring(0, 19)
|
||||
|
||||
console.log(`[Worker] ============================================================`)
|
||||
console.log(`[Worker] JOB DIMULAI — ID: ${job.id}`)
|
||||
console.log(`[Worker] Device : ${deviceId}`)
|
||||
console.log(`[Worker] Durasi : ${durationSeconds} detik`)
|
||||
console.log(`[Worker] Waktu : ${nowWIB} WIB (${nowUTC} UTC)`)
|
||||
console.log(`[Worker] Sumber : ${job.opts?.repeat?.cron ? 'Jadwal otomatis cron: ' + job.opts.repeat.cron : 'Siram manual'}`)
|
||||
console.log(`[Worker] ============================================================`)
|
||||
|
||||
// Langkah 1: Nyalakan pompa dengan kirim perintah ON ke relay ESP32
|
||||
publishRelay(deviceId, 'ON')
|
||||
console.log(`[Worker] >> Relay ON dikirim ke ${deviceId}`)
|
||||
|
||||
// Kirim Notifikasi OneSignal
|
||||
try {
|
||||
const { data: device } = await supabase
|
||||
.from('devices')
|
||||
.select('claimed_by')
|
||||
.eq('device_id', deviceId)
|
||||
.single()
|
||||
|
||||
if (device && device.claimed_by) {
|
||||
const isAuto = job.opts?.repeat?.cron
|
||||
sendNotification(
|
||||
device.claimed_by,
|
||||
'Penyiraman Dimulai 💦',
|
||||
`Pompa menyala selama ${durationSeconds} detik via ${isAuto ? 'Jadwal Otomatis' : 'Sistem'}.`
|
||||
)
|
||||
}
|
||||
} catch (e) {
|
||||
console.error('[Worker] Gagal mengirim notifikasi penyiraman:', e.message)
|
||||
}
|
||||
|
||||
// Langkah 2: Tunggu selama durasi yang ditentukan
|
||||
// Ini tidak memblokir server karena berjalan di proses worker yang terpisah
|
||||
await new Promise(resolve => setTimeout(resolve, durationSeconds * 1000))
|
||||
|
||||
// Langkah 3: Matikan pompa setelah durasi selesai
|
||||
publishRelay(deviceId, 'OFF')
|
||||
console.log(`[Worker] >> Relay OFF dikirim ke ${deviceId}`)
|
||||
console.log(`[Worker] JOB SELESAI — ${deviceId} | durasi: ${durationSeconds}s`)
|
||||
|
||||
}, { connection })
|
||||
|
||||
// ── Event Handlers ──────────────────────────────────────────────
|
||||
|
||||
worker.on('completed', (job) => {
|
||||
console.log(`[Worker] Job ${job.id} berhasil diselesaikan.`)
|
||||
})
|
||||
|
||||
worker.on('failed', (job, err) => {
|
||||
console.error(`[Worker] Job ${job?.id} GAGAL:`, err.message)
|
||||
console.error(`[Worker] Stack:`, err.stack)
|
||||
})
|
||||
|
||||
worker.on('error', (err) => {
|
||||
console.error(`[Worker] Worker error:`, err.message)
|
||||
})
|
||||
|
||||
worker.on('ready', () => {
|
||||
console.log(`[Worker] Worker siap memproses antrian 'irrigation'`)
|
||||
})
|
||||
|
||||
module.exports = worker
|
||||
|
|
@ -0,0 +1,84 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* SCHEDULE RESTORE — Pemulihan Jadwal Saat Server Start
|
||||
* =============================================================================
|
||||
* Masalah yang dipecahkan:
|
||||
* BullMQ menyimpan repeatable job di Redis, tapi jika terjadi:
|
||||
* - Server restart (Railway redeploy)
|
||||
* - Migrasi Redis
|
||||
* - Format key berubah (misal upgrade BullMQ dari v4 ke v5)
|
||||
* ...maka job di BullMQ bisa hilang atau tidak sinkron dengan database.
|
||||
*
|
||||
* Solusinya: Setiap kali server start, kita:
|
||||
* 1. Hapus SEMUA repeatable job lama dari BullMQ (bersih-bersih)
|
||||
* 2. Baca semua jadwal is_active=true dari Supabase
|
||||
* 3. Daftarkan ulang ke BullMQ dengan format yang benar (BullMQ v5+: key 'cron')
|
||||
*
|
||||
* Ini menjamin sinkronisasi antara DB dan BullMQ setiap kali server hidup.
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
const supabase = require('../supabase/client')
|
||||
const { irrigationQueue } = require('./irrigationQueue')
|
||||
|
||||
/**
|
||||
* Restore semua jadwal aktif dari Supabase ke BullMQ.
|
||||
* Dipanggil satu kali saat server startup di index.js.
|
||||
*/
|
||||
async function restoreSchedules() {
|
||||
console.log('[Restore] Memulai pemulihan jadwal dari database...')
|
||||
|
||||
try {
|
||||
// Langkah 1: Bersihkan SEMUA repeatable job lama dari BullMQ
|
||||
// Ini menghapus job yang mungkin tersimpan dengan format lama (pattern)
|
||||
const existingRepeatables = await irrigationQueue.getRepeatableJobs()
|
||||
console.log(`[Restore] Ditemukan ${existingRepeatables.length} job lama di BullMQ, membersihkan...`)
|
||||
|
||||
for (const job of existingRepeatables) {
|
||||
await irrigationQueue.removeRepeatableByKey(job.key)
|
||||
}
|
||||
console.log('[Restore] Job lama berhasil dihapus.')
|
||||
|
||||
// Langkah 2: Ambil semua jadwal aktif dari Supabase
|
||||
const { data: schedules, error } = await supabase
|
||||
.from('schedules')
|
||||
.select('*')
|
||||
.eq('is_active', true)
|
||||
|
||||
if (error) {
|
||||
console.error('[Restore] Gagal mengambil jadwal dari DB:', error.message)
|
||||
return
|
||||
}
|
||||
|
||||
if (!schedules || schedules.length === 0) {
|
||||
console.log('[Restore] Tidak ada jadwal aktif untuk di-restore.')
|
||||
return
|
||||
}
|
||||
|
||||
console.log(`[Restore] Mendaftarkan ulang ${schedules.length} jadwal aktif...`)
|
||||
|
||||
// Langkah 3: Daftarkan ulang setiap jadwal ke BullMQ dengan format BullMQ v5+
|
||||
for (const schedule of schedules) {
|
||||
try {
|
||||
const job = await irrigationQueue.add(
|
||||
'irrigate',
|
||||
{ deviceId: schedule.device_id, durationSeconds: schedule.duration_s },
|
||||
{
|
||||
repeat: { cron: schedule.cron }, // BullMQ v5+: key 'cron'
|
||||
jobId: schedule.bull_job_id,
|
||||
}
|
||||
)
|
||||
console.log(`[Restore] ✓ Jadwal ${schedule.id} (${schedule.cron}) → Job ${job.id}`)
|
||||
} catch (jobErr) {
|
||||
console.error(`[Restore] ✗ Gagal restore jadwal ${schedule.id}:`, jobErr.message)
|
||||
}
|
||||
}
|
||||
|
||||
console.log('[Restore] Pemulihan jadwal selesai.')
|
||||
} catch (err) {
|
||||
// Jangan sampai error di sini menghentikan server dari start
|
||||
console.error('[Restore] Error tidak terduga saat restore jadwal:', err.message)
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = { restoreSchedules }
|
||||
|
|
@ -0,0 +1,146 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* ROUTE: DEVICE — Klaim dan Manajemen Perangkat
|
||||
* =============================================================================
|
||||
* File ini mengurus proses "pairing" antara akun user dan perangkat ESP32 fisik.
|
||||
*
|
||||
* Kenapa perlu proses klaim?
|
||||
* Setiap ESP32 punya kode unik (claim_code) yang tercetak atau tertempel di box-nya.
|
||||
* Saat user pertama kali daftar, mereka harus memasukkan kode ini untuk
|
||||
* "mengklaim" perangkat tersebut sebagai milik mereka.
|
||||
* Setelah diklaim, semua data dan perintah akan dikaitkan dengan device itu.
|
||||
*
|
||||
* Endpoint yang tersedia:
|
||||
* - POST /api/device/claim → Klaim perangkat dengan kode unik
|
||||
* - GET /api/device/my-device/:userId → Ambil info perangkat milik user
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
const router = require('express').Router()
|
||||
const supabase = require('../supabase/client')
|
||||
|
||||
/**
|
||||
* POST /api/device/claim
|
||||
* -----------------------------------------------------------------------------
|
||||
* Menghubungkan perangkat ESP32 ke akun user berdasarkan kode klaim.
|
||||
*
|
||||
* Ada 4 pengecekan yang dilakukan secara berurutan:
|
||||
* 1. Validasi input — pastikan claim_code dan user_id dikirim
|
||||
* 2. Cari device — apakah kode ini terdaftar di database?
|
||||
* 3. Cek kepemilikan — sudah diklaim orang lain atau belum?
|
||||
* 4. Cek limit user — satu user hanya boleh punya satu device
|
||||
*
|
||||
* Body: { claim_code: string, user_id: string }
|
||||
*/
|
||||
router.post('/claim', async (req, res) => {
|
||||
const { claim_code, user_id } = req.body
|
||||
|
||||
// Pengecekan 1: Pastikan kedua field dikirim oleh client
|
||||
if (!claim_code || !user_id) {
|
||||
return res.status(400).json({ error: 'claim_code dan user_id wajib diisi' })
|
||||
}
|
||||
|
||||
// Pengecekan 2: Cari device dengan kode ini di database
|
||||
// claim_code selalu diubah ke huruf kapital agar tidak case-sensitive
|
||||
const { data: device, error } = await supabase
|
||||
.from('devices')
|
||||
.select('id, device_id, claimed_by')
|
||||
.eq('claim_code', claim_code.toUpperCase())
|
||||
.single()
|
||||
|
||||
if (error || !device) {
|
||||
return res.status(404).json({ error: 'Kode tidak ditemukan' })
|
||||
}
|
||||
|
||||
// Pengecekan 3: Apakah device ini sudah diklaim oleh orang LAIN?
|
||||
// Kalau yang mengklaim adalah user yang sama (misalnya retry), tetap diizinkan
|
||||
if (device.claimed_by && device.claimed_by !== user_id) {
|
||||
return res.status(409).json({ error: 'Device sudah diklaim oleh pengguna lain' })
|
||||
}
|
||||
|
||||
// Pengecekan 4: Apakah user ini sudah punya device lain sebelumnya?
|
||||
// Satu akun hanya boleh mengklaim satu perangkat
|
||||
const { data: existing } = await supabase
|
||||
.from('devices')
|
||||
.select('device_id')
|
||||
.eq('claimed_by', user_id)
|
||||
.single()
|
||||
|
||||
if (existing) {
|
||||
return res.status(409).json({ error: 'Anda sudah memiliki device terdaftar' })
|
||||
}
|
||||
|
||||
// Semua pengecekan lolos → Update database: tandai device sebagai milik user ini
|
||||
const { data: updated } = await supabase
|
||||
.from('devices')
|
||||
.update({ claimed_by: user_id, claimed_at: new Date() })
|
||||
.eq('claim_code', claim_code.toUpperCase())
|
||||
.select('device_id, label, location')
|
||||
.single()
|
||||
|
||||
// 5. Buat data threshold default otomatis untuk device ini
|
||||
// Mencegah error 500 saat user pertama kali mengatur batas sensor
|
||||
if (updated?.device_id) {
|
||||
const { error: thresholdErr } = await supabase
|
||||
.from('thresholds')
|
||||
.insert({
|
||||
device_id: updated.device_id,
|
||||
temp_max: 30,
|
||||
hum_max: 80,
|
||||
})
|
||||
if (thresholdErr) {
|
||||
console.warn('[Claim] Gagal membuat threshold default:', thresholdErr.message)
|
||||
} else {
|
||||
console.log(`[Claim] Threshold default dibuat untuk device: ${updated.device_id}`)
|
||||
}
|
||||
}
|
||||
|
||||
// Kembalikan info device yang berhasil diklaim ke client (Flutter)
|
||||
res.json({ message: 'Device berhasil diklaim', device: updated })
|
||||
})
|
||||
|
||||
/**
|
||||
* GET /api/device/my-device/:userId
|
||||
* -----------------------------------------------------------------------------
|
||||
* Mengambil informasi perangkat yang dimiliki oleh user tertentu.
|
||||
* Biasanya dipanggil saat user login untuk mengetahui device_id mereka,
|
||||
* lalu device_id itu digunakan di semua request API selanjutnya.
|
||||
*
|
||||
* Params: userId — UUID user dari Supabase Auth
|
||||
*/
|
||||
router.get('/my-device/:userId', async (req, res) => {
|
||||
console.log(`[Device Route] Handler berjalan untuk user: ${req.params.userId}`)
|
||||
const { data, error } = await supabase
|
||||
.from('devices')
|
||||
.select('device_id, label, location, is_online, last_seen')
|
||||
.eq('claimed_by', req.params.userId)
|
||||
.single()
|
||||
|
||||
if (error) {
|
||||
console.error('[Device Router] Gagal ambil my-device:', error.message)
|
||||
return res.status(404).json({ error: 'Device tidak ditemukan' })
|
||||
}
|
||||
res.json(data)
|
||||
})
|
||||
|
||||
/**
|
||||
* GET /api/device/status/:deviceId
|
||||
* -----------------------------------------------------------------------------
|
||||
* Mengambil status keaktifan perangkat (is_online, last_seen) secara real-time.
|
||||
* Params: deviceId — ID perangkat ESP32
|
||||
*/
|
||||
router.get('/status/:deviceId', async (req, res) => {
|
||||
const { data, error } = await supabase
|
||||
.from('devices')
|
||||
.select('is_online, last_seen')
|
||||
.eq('device_id', req.params.deviceId)
|
||||
.maybeSingle()
|
||||
|
||||
if (error || !data) {
|
||||
console.error('[Device Router] Gagal ambil status:', error ? error.message : 'Not found')
|
||||
return res.status(404).json({ error: 'Device tidak ditemukan' })
|
||||
}
|
||||
res.json(data)
|
||||
})
|
||||
|
||||
module.exports = router
|
||||
|
|
@ -0,0 +1,109 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* ROUTE: HISTORY — Riwayat Data Sensor
|
||||
* =============================================================================
|
||||
* File ini mengurus pengambilan data historis dari sensor SHT31.
|
||||
* Data historis berguna untuk:
|
||||
* - Menampilkan grafik suhu & kelembapan di aplikasi Flutter
|
||||
* - Melihat tren kondisi kumbung jamur dalam beberapa hari terakhir
|
||||
* - Evaluasi apakah sistem otomasi bekerja dengan baik
|
||||
*
|
||||
* Endpoint yang tersedia:
|
||||
* - GET /api/history/:deviceId → Data sensor terbaru (N data terakhir)
|
||||
* - GET /api/history/:deviceId/daily → Rata-rata per hari (N hari terakhir)
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
const router = require('express').Router()
|
||||
const supabase = require('../supabase/client')
|
||||
|
||||
/**
|
||||
* GET /api/history/:deviceId?limit=100
|
||||
* -----------------------------------------------------------------------------
|
||||
* Mengambil data log sensor terbaru untuk sebuah device.
|
||||
* Data diurutkan dari yang paling baru ke paling lama.
|
||||
*
|
||||
* Kenapa ada batas maksimum 500?
|
||||
* Untuk mencegah query yang terlalu berat ke database, yang bisa memperlambat
|
||||
* server dan menghabiskan kuota Supabase. Kalau butuh lebih banyak data,
|
||||
* sebaiknya gunakan endpoint /daily dengan agregasi.
|
||||
*
|
||||
* Params: deviceId — ID perangkat
|
||||
* Query: limit (opsional, default 100, maks 500)
|
||||
*
|
||||
* Contoh: GET /api/history/esp32-01?limit=50
|
||||
*/
|
||||
router.get('/:deviceId', async (req, res) => {
|
||||
console.log(`[History Route] Handler berjalan untuk device: ${req.params.deviceId}`)
|
||||
// Math.min memastikan limit tidak melebihi 500, sekali pun user kirim nilai lebih besar
|
||||
const limit = Math.min(parseInt(req.query.limit) || 100, 500)
|
||||
|
||||
const { data, error } = await supabase
|
||||
.from('sensor_logs')
|
||||
.select('temperature, humidity, relay_state, mode, created_at')
|
||||
.eq('device_id', req.params.deviceId)
|
||||
.order('created_at', { ascending: false }) // Data terbaru di atas
|
||||
.limit(limit)
|
||||
|
||||
if (error) {
|
||||
console.error('[History Router] Gagal ambil history:', error.message)
|
||||
return res.status(500).json({ error: error.message })
|
||||
}
|
||||
res.json(data)
|
||||
})
|
||||
|
||||
/**
|
||||
* GET /api/history/:deviceId/daily?days=7
|
||||
* -----------------------------------------------------------------------------
|
||||
* Mengambil data rata-rata suhu dan kelembapan per hari.
|
||||
* Berguna untuk grafik tren jangka panjang tanpa data yang terlalu padat.
|
||||
*
|
||||
* Endpoint ini memanggil Stored Procedure (fungsi SQL) bernama 'get_daily_average'
|
||||
* yang sudah didefinisikan di Supabase. Artinya proses agregasi dilakukan
|
||||
* di sisi database (lebih efisien daripada di server Node.js).
|
||||
*
|
||||
* Params: deviceId — ID perangkat
|
||||
* Query: days (opsional, default 7 hari)
|
||||
*
|
||||
* Contoh: GET /api/history/esp32-01/daily?days=14
|
||||
*/
|
||||
router.get('/:deviceId/daily', async (req, res) => {
|
||||
const { deviceId } = req.params
|
||||
const days = parseInt(req.query.days) || 7 // default 7 hari terakhir
|
||||
|
||||
// Panggil stored procedure di Supabase dengan parameter deviceId dan jumlah hari
|
||||
const { data, error } = await supabase.rpc('get_daily_average', {
|
||||
p_device_id: deviceId,
|
||||
p_days: days,
|
||||
})
|
||||
|
||||
if (error) return res.status(500).json({ error: error.message })
|
||||
res.json(data)
|
||||
})
|
||||
|
||||
/**
|
||||
* GET /api/history/:deviceId/hourly?days=7
|
||||
* -----------------------------------------------------------------------------
|
||||
* Mengambil rata-rata suhu dan kelembapan per jam.
|
||||
* Digunakan untuk grafik harian yang menunjukkan data per jam di UI.
|
||||
* Memanggil fungsi RPC 'get_hourly_average' di Supabase untuk performa yang lebih baik,
|
||||
* terutama mengingat data sensor masuk sangat cepat (misal setiap 3 detik).
|
||||
*
|
||||
* Params: deviceId — ID perangkat
|
||||
* Query: days (opsional, default 7 hari)
|
||||
*/
|
||||
router.get('/:deviceId/hourly', async (req, res) => {
|
||||
const { deviceId } = req.params
|
||||
const days = parseInt(req.query.days) || 7
|
||||
|
||||
// Panggil stored procedure di Supabase dengan parameter deviceId dan jumlah hari
|
||||
const { data, error } = await supabase.rpc('get_hourly_average', {
|
||||
p_device_id: deviceId,
|
||||
p_days: days,
|
||||
})
|
||||
|
||||
if (error) return res.status(500).json({ error: error.message })
|
||||
res.json(data)
|
||||
})
|
||||
|
||||
module.exports = router
|
||||
|
|
@ -0,0 +1,28 @@
|
|||
const router = require('express').Router()
|
||||
const { measureLatency } = require('../mqtt/mqttClient')
|
||||
|
||||
/**
|
||||
* GET /api/latency
|
||||
* -----------------------------------------------------------------------------
|
||||
* Mengukur RTT (Round-Trip Time) dari backend ke MQTT Broker (HiveMQ Cloud).
|
||||
*/
|
||||
router.get('/', async (req, res) => {
|
||||
try {
|
||||
const rtt = await measureLatency()
|
||||
const oneWay = Math.round(rtt / 2)
|
||||
res.json({
|
||||
success: true,
|
||||
backend_to_broker_rtt_ms: rtt,
|
||||
estimated_one_way_latency_ms: oneWay,
|
||||
message: `Pengukuran berhasil. RTT: ${rtt} ms.`
|
||||
})
|
||||
} catch (err) {
|
||||
console.error('[Latency Route] Gagal mengukur latensi:', err.message)
|
||||
res.status(500).json({
|
||||
success: false,
|
||||
error: err.message
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
module.exports = router
|
||||
|
|
@ -0,0 +1,151 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* ROUTE: MODE — Kontrol Mode Operasi ESP32
|
||||
* =============================================================================
|
||||
* File ini mengurus pergantian mode operasi perangkat ESP32 secara remote.
|
||||
*
|
||||
* ESP32 mendukung 3 mode:
|
||||
* - auto : Pompa menyala/mati otomatis berdasarkan pembacaan sensor SHT31
|
||||
* - manual : Pompa dikendalikan sepenuhnya oleh perintah dari aplikasi
|
||||
* - offline : Mode darurat saat koneksi internet terputus (logika auto, tanpa MQTT publish)
|
||||
*
|
||||
* Cara Kerja:
|
||||
* 1. Aplikasi Flutter memanggil POST /api/mode/:deviceId
|
||||
* 2. Backend memvalidasi nilai mode
|
||||
* 3. Backend mengirim perintah ke ESP32 via MQTT (cmd/mode/<deviceId>)
|
||||
* 4. Backend menyimpan mode baru ke kolom current_mode di tabel devices
|
||||
* 5. Backend mengirim notifikasi ke user bahwa mode telah berubah
|
||||
*
|
||||
* Endpoint yang tersedia:
|
||||
* - GET /api/mode/:deviceId → Baca mode aktif saat ini dari database
|
||||
* - POST /api/mode/:deviceId → Ubah mode operasi ESP32
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
const router = require('express').Router()
|
||||
const supabase = require('../supabase/client')
|
||||
const { publishMode } = require('../mqtt/mqttClient')
|
||||
const { sendNotification } = require('../utils/notification')
|
||||
|
||||
// Daftar nilai mode yang diizinkan (whitelist)
|
||||
const VALID_MODES = ['auto', 'manual', 'offline']
|
||||
|
||||
// Label ramah untuk notifikasi
|
||||
const MODE_LABELS = {
|
||||
auto: 'AUTO 🤖 — Pompa dikendalikan sensor secara otomatis',
|
||||
manual: 'MANUAL 🖐️ — Pompa dikendalikan langsung dari aplikasi',
|
||||
offline: 'OFFLINE 📴 — Mode darurat tanpa koneksi cloud',
|
||||
}
|
||||
|
||||
/**
|
||||
* GET /api/mode/:deviceId
|
||||
* -----------------------------------------------------------------------------
|
||||
* Membaca mode operasi aktif dari database untuk sebuah device.
|
||||
* Biasanya dipanggil saat aplikasi Flutter dibuka untuk menampilkan
|
||||
* status mode yang sedang berjalan.
|
||||
*
|
||||
* Params: deviceId — ID perangkat, contoh: "esp32-new-05"
|
||||
*/
|
||||
router.get('/:deviceId', async (req, res) => {
|
||||
const { data, error } = await supabase
|
||||
.from('devices')
|
||||
.select('device_id, current_mode, is_online, last_seen')
|
||||
.eq('device_id', req.params.deviceId)
|
||||
.single()
|
||||
|
||||
if (error || !data) {
|
||||
return res.status(404).json({ error: 'Device tidak ditemukan' })
|
||||
}
|
||||
|
||||
res.json({
|
||||
device_id: data.device_id,
|
||||
current_mode: data.current_mode ?? 'auto', // Fallback ke 'auto' jika kolom masih null
|
||||
is_online: data.is_online,
|
||||
last_seen: data.last_seen,
|
||||
})
|
||||
})
|
||||
|
||||
/**
|
||||
* POST /api/mode/:deviceId
|
||||
* -----------------------------------------------------------------------------
|
||||
* Mengubah mode operasi ESP32 dan menyimpan perubahan ke database.
|
||||
*
|
||||
* Alur eksekusi:
|
||||
* 1. Validasi nilai mode (hanya "auto", "manual", "offline" yang diterima)
|
||||
* 2. Cek apakah device terdaftar dan dimiliki oleh seseorang
|
||||
* 3. Kirim perintah mode baru ke ESP32 via MQTT
|
||||
* 4. Update kolom current_mode di tabel devices di Supabase
|
||||
* 5. Kirim notifikasi OneSignal ke pemilik device
|
||||
*
|
||||
* Params: deviceId — ID perangkat
|
||||
* Body: { mode: "auto" | "manual" | "offline" }
|
||||
*/
|
||||
router.post('/:deviceId', async (req, res) => {
|
||||
const { deviceId } = req.params
|
||||
const { mode } = req.body
|
||||
|
||||
// Validasi 1: Field mode wajib dikirim
|
||||
if (!mode) {
|
||||
return res.status(400).json({ error: 'Field "mode" wajib diisi' })
|
||||
}
|
||||
|
||||
// Validasi 2: Nilai mode harus salah satu dari whitelist
|
||||
if (!VALID_MODES.includes(mode)) {
|
||||
return res.status(400).json({
|
||||
error: `Mode tidak valid. Nilai yang diizinkan: ${VALID_MODES.join(', ')}`,
|
||||
})
|
||||
}
|
||||
|
||||
// Ambil info device (dibutuhkan untuk validasi keberadaan + kirim notifikasi)
|
||||
const { data: device, error: deviceErr } = await supabase
|
||||
.from('devices')
|
||||
.select('device_id, claimed_by, current_mode')
|
||||
.eq('device_id', deviceId)
|
||||
.single()
|
||||
|
||||
if (deviceErr || !device) {
|
||||
return res.status(404).json({ error: 'Device tidak ditemukan' })
|
||||
}
|
||||
|
||||
// Kirim perintah ganti mode ke ESP32 via MQTT
|
||||
// Selalu kirim, karena ESP32 mungkin sebelumnya sedang "cooldown" dan mengabaikan perintah
|
||||
publishMode(deviceId, mode)
|
||||
|
||||
// Simpan mode baru ke database hanya jika berbeda
|
||||
if (device.current_mode !== mode) {
|
||||
const { error: updateErr } = await supabase
|
||||
.from('devices')
|
||||
.update({ current_mode: mode })
|
||||
.eq('device_id', deviceId)
|
||||
|
||||
if (updateErr) {
|
||||
return res.status(500).json({ error: 'Gagal menyimpan mode ke database' })
|
||||
}
|
||||
}
|
||||
|
||||
// Kirim notifikasi ke pemilik device hanya jika mode benar-benar berubah
|
||||
if (device.claimed_by && device.current_mode !== mode) {
|
||||
try {
|
||||
sendNotification(
|
||||
device.claimed_by,
|
||||
'Mode Kumbung Berubah 🔄',
|
||||
`Mode diubah ke ${MODE_LABELS[mode]}`
|
||||
)
|
||||
} catch (notifErr) {
|
||||
// Jangan sampai gagal notifikasi membatalkan seluruh respons
|
||||
console.error('[Mode] Gagal kirim notifikasi:', notifErr.message)
|
||||
}
|
||||
}
|
||||
|
||||
console.log(`[Mode] ${deviceId} → ${device.current_mode} → ${mode} (Forced MQTT Publish)`)
|
||||
|
||||
res.json({
|
||||
message: device.current_mode === mode
|
||||
? `Perintah mode ${mode} dikirim ulang ke perangkat`
|
||||
: `Mode berhasil diubah ke ${mode}`,
|
||||
mode,
|
||||
changed: device.current_mode !== mode,
|
||||
})
|
||||
})
|
||||
|
||||
module.exports = router
|
||||
|
|
@ -0,0 +1,264 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* ROUTE: SCHEDULE — Jadwal & Kontrol Penyiraman
|
||||
* =============================================================================
|
||||
* File ini adalah "pusat kontrol" untuk semua hal yang berhubungan dengan
|
||||
* penyiraman: jadwal terjadwal, trigger manual, hingga menghentikan pompa.
|
||||
*
|
||||
* Konsep penting: Hybrid antara Database dan BullMQ
|
||||
* Semua jadwal disimpan di DUA tempat sekaligus:
|
||||
* 1. Database Supabase → untuk menyimpan data permanen (label, cron, durasi, dll)
|
||||
* 2. BullMQ (Redis) → untuk mengeksekusi jadwal tepat waktu
|
||||
*
|
||||
* Kenapa dua tempat? Karena BullMQ menjalankan jadwal berdasarkan waktu nyata,
|
||||
* tapi jika server di-restart, semua data BullMQ hilang. Dengan menyimpan di
|
||||
* Supabase, jadwal bisa di-restore saat server dinyalakan kembali.
|
||||
*
|
||||
* Endpoint yang tersedia:
|
||||
* - GET /api/schedule/:deviceId → Daftar semua jadwal
|
||||
* - POST /api/schedule/:deviceId → Buat jadwal baru
|
||||
* - DELETE /api/schedule/:id → Hapus jadwal
|
||||
* - POST /api/schedule/:deviceId/now → Siram sekarang (sekali jalan)
|
||||
* - POST /api/schedule/:deviceId/stop → Hentikan pompa sekarang
|
||||
* - PATCH /api/schedule/:id/toggle → Aktifkan atau nonaktifkan jadwal
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
const router = require('express').Router()
|
||||
const supabase = require('../supabase/client')
|
||||
const { irrigationQueue } = require('../queues/irrigationQueue')
|
||||
const { publishRelay } = require('../mqtt/mqttClient')
|
||||
|
||||
/**
|
||||
* GET /api/schedule/:deviceId
|
||||
* -----------------------------------------------------------------------------
|
||||
* Mengambil semua jadwal penyiraman yang dimiliki oleh sebuah device.
|
||||
* Diurutkan berdasarkan waktu pembuatan (paling lama di atas).
|
||||
* Biasanya dipanggil saat aplikasi Flutter membuka halaman "Jadwal".
|
||||
*
|
||||
* Params: deviceId — ID perangkat
|
||||
*/
|
||||
router.get('/:deviceId', async (req, res) => {
|
||||
const { data, error } = await supabase
|
||||
.from('schedules')
|
||||
.select('*')
|
||||
.eq('device_id', req.params.deviceId)
|
||||
.order('created_at')
|
||||
|
||||
if (error) return res.status(500).json({ error: error.message })
|
||||
res.json(data)
|
||||
})
|
||||
|
||||
/**
|
||||
* POST /api/schedule/:deviceId
|
||||
* -----------------------------------------------------------------------------
|
||||
* Membuat jadwal penyiraman baru yang akan berulang sesuai pola cron.
|
||||
*
|
||||
* Yang terjadi saat endpoint ini dipanggil:
|
||||
* 1. Validasi input (cron dan duration_s wajib ada)
|
||||
* 2. Daftarkan job repeatable ke BullMQ (yang akan mengeksekusinya nanti)
|
||||
* 3. Simpan data jadwal ke Supabase beserta ID job dari BullMQ
|
||||
*
|
||||
* Format Cron: "menit jam hari bulan hari-minggu"
|
||||
* Contoh:
|
||||
* - "0 6 * * *" → Setiap hari jam 06:00
|
||||
* - "30 17 * * *" → Setiap hari jam 17:30
|
||||
* - "0 6,17 * * *" → Setiap hari jam 06:00 dan 17:00
|
||||
*
|
||||
* Params: deviceId — ID perangkat
|
||||
* Body: { label: string, cron: string, duration_s: number }
|
||||
*/
|
||||
router.post('/:deviceId', async (req, res) => {
|
||||
const { label, cron, duration_s } = req.body
|
||||
const deviceId = req.params.deviceId
|
||||
|
||||
// Validasi: cron pattern dan durasi penyiraman wajib dikirim
|
||||
if (!cron || !duration_s) {
|
||||
return res.status(400).json({ error: 'cron dan duration_s wajib diisi' })
|
||||
}
|
||||
|
||||
// Validasi format cron (minimal harus ada 5 bagian yang dipisahkan spasi)
|
||||
if (cron.split(' ').length !== 5) {
|
||||
return res.status(400).json({ error: 'Format cron tidak valid' })
|
||||
}
|
||||
|
||||
// Langkah 1: Daftarkan job berulang ke BullMQ menggunakan pola cron
|
||||
// PENTING: BullMQ v5+ menggunakan key 'cron' bukan 'pattern' (breaking change)
|
||||
const jobId = `schedule-${deviceId}-${Date.now()}`
|
||||
const job = await irrigationQueue.add(
|
||||
'irrigate',
|
||||
{ deviceId, durationSeconds: duration_s },
|
||||
{
|
||||
repeat: { cron }, // BullMQ v5+: gunakan 'cron', bukan 'pattern'
|
||||
jobId, // ID unik untuk tiap jadwal
|
||||
}
|
||||
)
|
||||
|
||||
// Langkah 2: Simpan jadwal ke database beserta bull_job_id untuk referensi nanti
|
||||
// bull_job_id dibutuhkan saat ingin menghapus atau toggle jadwal
|
||||
const { data, error } = await supabase
|
||||
.from('schedules')
|
||||
.insert({ device_id: deviceId, label, cron, duration_s, bull_job_id: job.id })
|
||||
.select()
|
||||
.single()
|
||||
|
||||
if (error) return res.status(500).json({ error: error.message })
|
||||
res.status(201).json(data)
|
||||
})
|
||||
|
||||
/**
|
||||
* DELETE /api/schedule/:id
|
||||
* -----------------------------------------------------------------------------
|
||||
* Menghapus jadwal secara permanen dari database DAN dari antrian BullMQ.
|
||||
*
|
||||
* Penting: Harus hapus dari keduanya! Kalau hanya hapus dari database tapi
|
||||
* tidak dari BullMQ, job tetap akan jalan di waktu yang sudah dijadwalkan.
|
||||
*
|
||||
* Params: id — UUID jadwal dari tabel schedules di Supabase
|
||||
*/
|
||||
router.delete('/:id', async (req, res) => {
|
||||
// Langkah 1: Ambil detail jadwal dulu untuk mendapat bull_job_id dan cron pattern
|
||||
// Kedua info ini dibutuhkan untuk menghapus job dari BullMQ
|
||||
const { data: schedule, error: fetchErr } = await supabase
|
||||
.from('schedules')
|
||||
.select('bull_job_id, cron, device_id')
|
||||
.eq('id', req.params.id)
|
||||
.single()
|
||||
|
||||
if (fetchErr) return res.status(404).json({ error: 'Jadwal tidak ditemukan' })
|
||||
|
||||
// Langkah 2: Hapus job dari BullMQ agar tidak dieksekusi lagi
|
||||
// BullMQ v5+: gunakan 'cron', bukan 'pattern'
|
||||
await irrigationQueue.removeRepeatable('irrigate', {
|
||||
cron: schedule.cron,
|
||||
jobId: schedule.bull_job_id,
|
||||
})
|
||||
|
||||
// Langkah 3: Hapus record dari database Supabase
|
||||
await supabase.from('schedules').delete().eq('id', req.params.id)
|
||||
|
||||
res.json({ message: 'Jadwal dihapus' })
|
||||
})
|
||||
|
||||
/**
|
||||
* POST /api/schedule/:deviceId/now
|
||||
* -----------------------------------------------------------------------------
|
||||
* Memicu penyiraman SEKARANG JUGA tanpa membuat jadwal permanen.
|
||||
* Cocok untuk penyiraman darurat atau percobaan manual.
|
||||
*
|
||||
* Job ini tidak akan berulang (tidak pakai cron), langsung diproses Worker
|
||||
* segera setelah request masuk.
|
||||
*
|
||||
* Endpoint ini dilindungi rate limiter ketat (5x/menit) yang didefinisikan
|
||||
* di index.js untuk mencegah penyiraman berlebihan.
|
||||
*
|
||||
* Params: deviceId — ID perangkat
|
||||
* Body: { duration_s: number } — opsional, default 30 detik
|
||||
*/
|
||||
router.post('/:deviceId/now', async (req, res) => {
|
||||
const { duration_s = 30 } = req.body // Default durasi 30 detik jika tidak dikirim
|
||||
|
||||
// Tambahkan job satu kali ke antrian (tanpa repeat) dengan delay 0 (langsung)
|
||||
await irrigationQueue.add(
|
||||
'irrigate',
|
||||
{ deviceId: req.params.deviceId, durationSeconds: duration_s },
|
||||
{ delay: 0 } // langsung diproses tanpa jeda
|
||||
)
|
||||
|
||||
res.json({ message: `Siram manual ${duration_s}s dijadwalkan` })
|
||||
})
|
||||
|
||||
/**
|
||||
* POST /api/schedule/:deviceId/stop
|
||||
* -----------------------------------------------------------------------------
|
||||
* Menghentikan pompa secara paksa dan instan via MQTT.
|
||||
* Berguna jika user ingin menghentikan penyiraman sebelum waktunya habis.
|
||||
*
|
||||
* Catatan: Ini mengirim perintah OFF langsung ke relay, bukan membatalkan job
|
||||
* di BullMQ. Jadi jika ada job yang sedang berjalan, Worker masih akan
|
||||
* menunggu sampai durasi habis baru kirim OFF lagi (tidak berbahaya, hanya redundan).
|
||||
*
|
||||
* Params: deviceId — ID perangkat
|
||||
*/
|
||||
router.post('/:deviceId/stop', async (req, res) => {
|
||||
const { deviceId } = req.params
|
||||
|
||||
// Kirim perintah OFF langsung ke relay ESP32 via MQTT
|
||||
publishRelay(deviceId, 'OFF')
|
||||
|
||||
// Catat event ini ke tabel sensor_logs sebagai audit trail
|
||||
await supabase.from('sensor_logs').insert({
|
||||
device_id: deviceId,
|
||||
event: 'manual_stop',
|
||||
note: 'Pompa dimatikan manual',
|
||||
})
|
||||
|
||||
res.json({ message: 'Pompa dimatikan' })
|
||||
})
|
||||
|
||||
/**
|
||||
* PATCH /api/schedule/:id/toggle
|
||||
* -----------------------------------------------------------------------------
|
||||
* Mengaktifkan atau menonaktifkan jadwal tanpa menghapusnya.
|
||||
* Ini lebih nyaman bagi user daripada harus delete lalu create ulang.
|
||||
*
|
||||
* Cara kerjanya:
|
||||
* - Jika jadwal sedang AKTIF → hapus job dari BullMQ, update is_active = false
|
||||
* - Jika jadwal sedang NON-AKTIF → daftarkan ulang job ke BullMQ, update is_active = true
|
||||
*
|
||||
* Data di database tetap tersimpan dalam kedua kondisi, hanya status is_active yang berubah.
|
||||
*
|
||||
* Params: id — UUID jadwal dari tabel schedules
|
||||
*/
|
||||
router.patch('/:id/toggle', async (req, res) => {
|
||||
// Ambil data jadwal saat ini untuk mengetahui status is_active dan detail cron-nya
|
||||
const { data: schedule, error: fetchErr } = await supabase
|
||||
.from('schedules')
|
||||
.select('is_active, cron, duration_s, device_id, bull_job_id')
|
||||
.eq('id', req.params.id)
|
||||
.single()
|
||||
|
||||
if (fetchErr) return res.status(404).json({ error: 'Jadwal tidak ditemukan' })
|
||||
|
||||
// Status baru adalah kebalikan dari status saat ini
|
||||
const newStatus = !schedule.is_active
|
||||
|
||||
if (newStatus) {
|
||||
// Jadwal akan DIAKTIFKAN: daftarkan kembali job ke BullMQ
|
||||
// BullMQ v5+: gunakan 'cron', bukan 'pattern'
|
||||
const job = await irrigationQueue.add(
|
||||
'irrigate',
|
||||
{ deviceId: schedule.device_id, durationSeconds: schedule.duration_s },
|
||||
{
|
||||
repeat: { cron: schedule.cron }, // BullMQ v5+: key 'cron'
|
||||
jobId: schedule.bull_job_id, // gunakan ID yang sama agar referensi tidak berubah
|
||||
}
|
||||
)
|
||||
console.log(`[Schedule] Job ${job.id} diaktifkan kembali`)
|
||||
} else {
|
||||
// Jadwal akan DINONAKTIFKAN: hapus job dari BullMQ (tapi data DB tetap ada)
|
||||
// BullMQ v5+: gunakan 'cron', bukan 'pattern'
|
||||
await irrigationQueue.removeRepeatable('irrigate', {
|
||||
cron: schedule.cron,
|
||||
jobId: schedule.bull_job_id,
|
||||
})
|
||||
console.log(`[Schedule] Job ${schedule.bull_job_id} dinonaktifkan`)
|
||||
}
|
||||
|
||||
// Update kolom is_active di database Supabase
|
||||
const { data, error } = await supabase
|
||||
.from('schedules')
|
||||
.update({ is_active: newStatus })
|
||||
.eq('id', req.params.id)
|
||||
.select()
|
||||
.single()
|
||||
|
||||
if (error) return res.status(500).json({ error: error.message })
|
||||
|
||||
res.json({
|
||||
message: newStatus ? 'Jadwal diaktifkan' : 'Jadwal dinonaktifkan',
|
||||
data,
|
||||
})
|
||||
})
|
||||
|
||||
module.exports = router
|
||||
|
|
@ -0,0 +1,93 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* ROUTE: THRESHOLD — Pengaturan Batas Suhu & Kelembapan
|
||||
* =============================================================================
|
||||
* File ini berfungsi mengatur "batas" yang menentukan kapan relay otomatis
|
||||
* menyala atau mati.
|
||||
*
|
||||
* Cara kerjanya:
|
||||
* Saat sensor SHT31 di ESP32 membaca data, backend akan membandingkan nilainya
|
||||
* dengan threshold yang tersimpan di sini:
|
||||
* - Jika suhu ATAU kelembapan melewati batas → relay ON
|
||||
* - Jika keduanya di bawah batas → relay OFF
|
||||
*
|
||||
* Endpoint yang tersedia:
|
||||
* - GET /api/threshold/:deviceId → Ambil threshold aktif untuk device tertentu
|
||||
* - POST /api/threshold/:deviceId → Update threshold (langsung efektif ke ESP32)
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
const router = require('express').Router()
|
||||
const supabase = require('../supabase/client')
|
||||
const { publishThreshold } = require('../mqtt/mqttClient')
|
||||
|
||||
/**
|
||||
* GET /api/threshold/:deviceId
|
||||
* -----------------------------------------------------------------------------
|
||||
* Mengambil nilai threshold yang sedang aktif untuk sebuah device.
|
||||
* Biasanya dipanggil saat aplikasi Flutter pertama kali dibuka untuk
|
||||
* menampilkan pengaturan saat ini kepada user.
|
||||
*
|
||||
* Params: deviceId — ID perangkat, contoh: "esp32-01"
|
||||
*/
|
||||
router.get('/:deviceId', async (req, res) => {
|
||||
const { data, error } = await supabase
|
||||
.from('thresholds')
|
||||
.select('*')
|
||||
.eq('device_id', req.params.deviceId)
|
||||
.single()
|
||||
|
||||
if (error) return res.status(404).json({ error: error.message })
|
||||
res.json(data)
|
||||
})
|
||||
|
||||
/**
|
||||
* POST /api/threshold/:deviceId
|
||||
* -----------------------------------------------------------------------------
|
||||
* Mengubah nilai threshold untuk sebuah device.
|
||||
*
|
||||
* Yang terjadi saat endpoint ini dipanggil:
|
||||
* 1. Validasi input (temp_max dan hum_max wajib ada)
|
||||
* 2. Update nilai di database Supabase
|
||||
* 3. Kirim nilai baru ke ESP32 via MQTT (langsung efektif!)
|
||||
* 4. Update cache in-memory agar tidak perlu tunggu 30 detik untuk efektif
|
||||
*
|
||||
* Jadi user tidak perlu reload atau menunggu — perubahan langsung terasa.
|
||||
*
|
||||
* Params: deviceId — ID perangkat
|
||||
* Body: { temp_max: number, hum_max: number }
|
||||
*/
|
||||
router.post('/:deviceId', async (req, res) => {
|
||||
const { temp_min, temp_max, hum_max } = req.body
|
||||
|
||||
// Validasi: kedua nilai threshold (temp_max & hum_max) wajib dikirim dan bukan null/undefined
|
||||
if (temp_max === undefined || temp_max === null || hum_max === undefined || hum_max === null) {
|
||||
return res.status(400).json({ error: 'temp_max dan hum_max wajib diisi' })
|
||||
}
|
||||
|
||||
// Jika temp_min tidak dikirim (misalnya request dari Flutter versi lama), beri nilai default 20.0
|
||||
const finalTempMin = (temp_min === undefined || temp_min === null) ? 20.0 : temp_min;
|
||||
|
||||
// Simpan nilai baru ke database
|
||||
const { data, error } = await supabase
|
||||
.from('thresholds')
|
||||
.update({
|
||||
temp_min: finalTempMin,
|
||||
temp_max,
|
||||
hum_max,
|
||||
updated_at: new Date()
|
||||
})
|
||||
.eq('device_id', req.params.deviceId)
|
||||
.select()
|
||||
.single()
|
||||
|
||||
if (error) return res.status(500).json({ error: error.message })
|
||||
|
||||
// Kirim nilai threshold baru langsung ke ESP32 via MQTT
|
||||
// Sekaligus memperbarui cache agar langsung efektif tanpa tunggu TTL
|
||||
publishThreshold(req.params.deviceId, finalTempMin, temp_max, hum_max)
|
||||
|
||||
res.json({ message: 'Threshold diupdate', data })
|
||||
})
|
||||
|
||||
module.exports = router
|
||||
|
|
@ -0,0 +1,27 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* SUPABASE CLIENT
|
||||
* =============================================================================
|
||||
* File ini punya satu tugas: bikin "jembatan" koneksi ke database Supabase.
|
||||
*
|
||||
* Kenapa pakai Service Key, bukan Anon Key?
|
||||
* - Anon Key dipakai di sisi client (Flutter/browser), aksesnya terbatas oleh
|
||||
* Row Level Security (RLS).
|
||||
* - Service Key dipakai di sisi backend (server ini), karena punya akses
|
||||
* penuh tanpa dibatasi RLS. Makanya JANGAN PERNAH bocorkan key ini ke publik!
|
||||
*
|
||||
* Cara pakainya: cukup import file ini dari mana saja di proyek ini.
|
||||
* Contoh: const supabase = require('../supabase/client')
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
const { createClient } = require('@supabase/supabase-js')
|
||||
|
||||
// Inisialisasi koneksi ke Supabase menggunakan URL dan Service Key
|
||||
// yang diambil dari file .env (supaya tidak hardcode di kode)
|
||||
const supabase = createClient(
|
||||
process.env.SUPABASE_URL,
|
||||
process.env.SUPABASE_SERVICE_KEY // pakai service key di backend, bukan anon key
|
||||
)
|
||||
|
||||
module.exports = supabase
|
||||
|
|
@ -0,0 +1,145 @@
|
|||
/**
|
||||
* =============================================================================
|
||||
* FAILOVER MANAGER — Pengendali Mode Backup & Utama
|
||||
* =============================================================================
|
||||
* Berfungsi untuk mengelola status server saat dijalankan di VPS (sebagai backup)
|
||||
* atau di Railway (sebagai primary).
|
||||
*
|
||||
* Cara kerjanya:
|
||||
* 1. Jika IS_BACKUP_SERVER !== 'true' (Primary / Railway):
|
||||
* - Langsung aktifkan semua service background (MQTT, Worker, Scheduler, Offline Detector).
|
||||
*
|
||||
* 2. Jika IS_BACKUP_SERVER === 'true' (Backup / VPS):
|
||||
* - Jangan aktifkan service background pada startup.
|
||||
* - Jalankan pemantauan berkala (ping) ke PRIMARY_SERVER_URL.
|
||||
* - Jika ping sukses (Primary hidup) -> Pastikan service background mati/nonaktif.
|
||||
* - Jika ping gagal/timeout (Primary mati) -> Aktifkan semua service background.
|
||||
* - Jika ping kembali sukses -> Matikan kembali service background.
|
||||
* =============================================================================
|
||||
*/
|
||||
|
||||
const { connect, disconnect } = require('../mqtt/mqttClient')
|
||||
const worker = require('../queues/irrigationWorker')
|
||||
const { restoreSchedules } = require('../queues/scheduleRestore')
|
||||
const { startOfflineDetector, stopOfflineDetector } = require('../jobs/offlineDetector')
|
||||
|
||||
const IS_BACKUP_SERVER = process.env.IS_BACKUP_SERVER === 'true'
|
||||
const PRIMARY_SERVER_URL = process.env.PRIMARY_SERVER_URL // Contoh: https://xxx.up.railway.app
|
||||
const PING_INTERVAL_MS = parseInt(process.env.FAILOVER_PING_INTERVAL_MS) || 10000 // default 10s
|
||||
|
||||
let isFailoverActive = false // Apakah backup server sedang melayani background tasks
|
||||
let pingIntervalId = null
|
||||
|
||||
async function pingPrimary() {
|
||||
if (!PRIMARY_SERVER_URL) {
|
||||
console.error('[Failover] PRIMARY_SERVER_URL tidak disetel! Tidak bisa memantau primary server.')
|
||||
return false
|
||||
}
|
||||
|
||||
let timeoutId = null
|
||||
try {
|
||||
const controller = new AbortController()
|
||||
timeoutId = setTimeout(() => controller.abort(), 5000) // 5s timeout
|
||||
|
||||
const response = await fetch(PRIMARY_SERVER_URL, {
|
||||
method: 'GET',
|
||||
signal: controller.signal
|
||||
})
|
||||
|
||||
return response.ok // true jika status 200-299
|
||||
} catch (err) {
|
||||
// Bisa error network atau timeout
|
||||
return false
|
||||
} finally {
|
||||
if (timeoutId) {
|
||||
clearTimeout(timeoutId)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function activateBackupServices() {
|
||||
if (isFailoverActive) return
|
||||
isFailoverActive = true
|
||||
console.log('[Failover] ⚠️ PRIMARY SERVER TERDETEKSI DOWN! Mengaktifkan layanan backup lokal...')
|
||||
|
||||
try {
|
||||
// 1. Hubungkan ke MQTT Broker
|
||||
connect()
|
||||
|
||||
// 2. Aktifkan Worker BullMQ
|
||||
await worker.resume()
|
||||
console.log('[Failover] Worker BullMQ di-resume.')
|
||||
|
||||
// 3. Restore Jadwal dari DB ke BullMQ
|
||||
await restoreSchedules()
|
||||
|
||||
// 4. Jalankan offline detector
|
||||
startOfflineDetector()
|
||||
|
||||
console.log('[Failover] ✅ Layanan backup berhasil diaktifkan dan siap melayani.')
|
||||
} catch (err) {
|
||||
console.error('[Failover] Gagal mengaktifkan beberapa layanan backup:', err.message)
|
||||
}
|
||||
}
|
||||
|
||||
async function deactivateBackupServices() {
|
||||
if (!isFailoverActive) return
|
||||
isFailoverActive = false
|
||||
console.log('[Failover] 💚 PRIMARY SERVER KEMBALI ONLINE! Menonaktifkan layanan backup lokal ke standby...')
|
||||
|
||||
try {
|
||||
// 1. Putuskan hubungan MQTT
|
||||
disconnect()
|
||||
|
||||
// 2. Pause Worker BullMQ
|
||||
await worker.pause()
|
||||
console.log('[Failover] Worker BullMQ di-pause.')
|
||||
|
||||
// 3. Matikan offline detector
|
||||
stopOfflineDetector()
|
||||
|
||||
console.log('[Failover] ✅ Layanan backup dinonaktifkan kembali ke mode standby.')
|
||||
} catch (err) {
|
||||
console.error('[Failover] Gagal menonaktifkan beberapa layanan backup:', err.message)
|
||||
}
|
||||
}
|
||||
|
||||
async function initFailover() {
|
||||
if (!IS_BACKUP_SERVER) {
|
||||
console.log('[Failover] Berjalan sebagai PRIMARY SERVER (Railway). Mengaktifkan semua service...')
|
||||
// Mode Primary: Langsung nyalakan semuanya
|
||||
connect()
|
||||
// BullMQ worker default-nya aktif saat dibuat, jadi tidak perlu resume.
|
||||
await restoreSchedules()
|
||||
startOfflineDetector()
|
||||
return
|
||||
}
|
||||
|
||||
console.log('[Failover] Berjalan sebagai BACKUP SERVER (VPS) dalam mode Standby.')
|
||||
console.log(`[Failover] Memantau primary server di: ${PRIMARY_SERVER_URL} setiap ${PING_INTERVAL_MS / 1000}s`)
|
||||
|
||||
// Mode Backup: Pause worker terlebih dahulu agar tidak memproses antrian saat standby
|
||||
try {
|
||||
await worker.pause()
|
||||
console.log('[Failover] Worker BullMQ di-pause (Standby).')
|
||||
} catch (err) {
|
||||
console.error('[Failover] Gagal mem-pause Worker BullMQ di startup:', err.message)
|
||||
}
|
||||
|
||||
// Jalankan monitoring loop
|
||||
pingIntervalId = setInterval(async () => {
|
||||
const isPrimaryAlive = await pingPrimary()
|
||||
|
||||
if (isPrimaryAlive) {
|
||||
if (isFailoverActive) {
|
||||
await deactivateBackupServices()
|
||||
}
|
||||
} else {
|
||||
if (!isFailoverActive) {
|
||||
await activateBackupServices()
|
||||
}
|
||||
}
|
||||
}, PING_INTERVAL_MS)
|
||||
}
|
||||
|
||||
module.exports = { initFailover }
|
||||
|
|
@ -0,0 +1,100 @@
|
|||
const Redis = require('ioredis');
|
||||
|
||||
const redis = process.env.REDIS_URL ? new Redis(process.env.REDIS_URL) : null;
|
||||
|
||||
// Fallback anti-spam in-memory jika Redis tidak tersedia
|
||||
const inMemoryCooldown = new Map();
|
||||
|
||||
/**
|
||||
* Kirim push notification via OneSignal dengan perlindungan anti-spam bawaan.
|
||||
* Setiap kombinasi (userId + title) diberi cooldown agar tidak bisa dikirim
|
||||
* berulang kali dalam jangka waktu tertentu.
|
||||
*
|
||||
* @param {string} userId - ID user penerima (external_id OneSignal)
|
||||
* @param {string} title - Judul notifikasi (juga dipakai sebagai kunci cooldown)
|
||||
* @param {string} message - Isi pesan notifikasi
|
||||
* @param {number} [cooldownSec=300] - Cooldown dalam detik (default: 5 menit)
|
||||
*/
|
||||
const sendNotification = async (userId, title, message, cooldownSec = 300) => {
|
||||
const appId = process.env.ONESIGNAL_APP_ID;
|
||||
const restApiKey = process.env.ONESIGNAL_REST_API_KEY;
|
||||
|
||||
if (!appId || !restApiKey) {
|
||||
console.warn('[Notif] OneSignal credentials missing. Skipping.');
|
||||
return;
|
||||
}
|
||||
|
||||
// Ensure targetIds is a flat array of unique strings
|
||||
const targetIds = Array.isArray(userId) ? [...new Set(userId.flat())] : [userId];
|
||||
|
||||
if (targetIds.length === 0) return;
|
||||
|
||||
// ── Anti-Spam Guard ──────────────────────────────────────────────────────
|
||||
// Buat kunci unik berdasarkan siapa penerima dan judul notifikasinya.
|
||||
// Ini mencegah notif "Penyiraman Dimulai" atau "Perangkat Offline"
|
||||
// dikirim berkali-kali dalam waktu singkat.
|
||||
const cooldownKey = `notif_cd:${targetIds.join(',')}:${title.replace(/\s+/g, '_').toLowerCase()}`;
|
||||
|
||||
if (redis) {
|
||||
const isOnCooldown = await redis.get(cooldownKey);
|
||||
if (isOnCooldown) {
|
||||
console.log(`[Notif] Cooldown aktif, skip: "${title}" → ${targetIds}`);
|
||||
return;
|
||||
}
|
||||
// Set cooldown key di Redis
|
||||
await redis.set(cooldownKey, '1', 'EX', cooldownSec);
|
||||
} else {
|
||||
// Fallback: gunakan in-memory Map jika Redis tidak ada
|
||||
const expiresAt = inMemoryCooldown.get(cooldownKey);
|
||||
if (expiresAt && Date.now() < expiresAt) {
|
||||
console.log(`[Notif] Cooldown aktif (in-memory), skip: "${title}" → ${targetIds}`);
|
||||
return;
|
||||
}
|
||||
inMemoryCooldown.set(cooldownKey, Date.now() + cooldownSec * 1000);
|
||||
}
|
||||
// ────────────────────────────────────────────────────────────────────────
|
||||
|
||||
const payload = {
|
||||
app_id: appId,
|
||||
include_aliases: {
|
||||
external_id: targetIds
|
||||
},
|
||||
target_channel: 'push',
|
||||
headings: { en: title },
|
||||
contents: { en: message },
|
||||
// -- Konfigurasi Grouping & Stacking (Agar rapi seperti WhatsApp) --
|
||||
android_group: 'jamur_monitoring_group', // Mengelompokkan notifikasi di Android
|
||||
thread_id: 'jamur_monitoring_group', // Mengelompokkan notifikasi di iOS
|
||||
// Menimpa (replace) notifikasi lama yang jenisnya sama, agar tidak spam
|
||||
collapse_id: title.replace(/\s+/g, '_').toLowerCase()
|
||||
};
|
||||
|
||||
try {
|
||||
const response = await fetch('https://onesignal.com/api/v1/notifications', {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
'Content-Type': 'application/json; charset=utf-8',
|
||||
'Authorization': `Basic ${restApiKey}`
|
||||
},
|
||||
body: JSON.stringify(payload)
|
||||
});
|
||||
|
||||
const result = await response.json();
|
||||
if (result.errors) {
|
||||
// Jika error karena user belum login (invalid_aliases), jangan tampilkan error panjang
|
||||
if (result.errors.invalid_aliases) {
|
||||
console.warn(`[Notif] User belum register OneSignal (invalid_aliases):`, targetIds);
|
||||
} else {
|
||||
console.error('[Notif] OneSignal error:', result.errors);
|
||||
}
|
||||
} else {
|
||||
console.log(`[Notif] Terkirim ke ${targetIds}: "${title}"`);
|
||||
}
|
||||
} catch (error) {
|
||||
console.error('[Notif] Gagal kirim via OneSignal:', error.message);
|
||||
}
|
||||
};
|
||||
|
||||
module.exports = {
|
||||
sendNotification
|
||||
};
|
||||
|
|
@ -0,0 +1,126 @@
|
|||
require('dotenv').config();
|
||||
const mqtt = require('mqtt');
|
||||
|
||||
console.log('=== HIVE MQTT LATENCY TESTER ===');
|
||||
console.log('Menghubungkan ke:', process.env.MQTT_BROKER_URL);
|
||||
|
||||
const client = mqtt.connect(process.env.MQTT_BROKER_URL, {
|
||||
port: parseInt(process.env.MQTT_PORT) || 8883,
|
||||
username: process.env.MQTT_USERNAME,
|
||||
password: process.env.MQTT_PASSWORD,
|
||||
protocol: 'mqtts',
|
||||
clientId: `latency-tester-${Date.now()}`,
|
||||
clean: true,
|
||||
});
|
||||
|
||||
const TOPIC_SENSOR = 'sensor/sht31';
|
||||
const TOPIC_RELAY_WILDCARD = 'cmd/relay/+';
|
||||
const TOPIC_PING = 'latency/test/ping';
|
||||
|
||||
// Map untuk mencatat waktu kirim perintah relay
|
||||
const pendingCommands = new Map();
|
||||
|
||||
// Map untuk melacak waktu pengiriman ping loopback
|
||||
const pingSentTimes = new Map();
|
||||
const pingPubAckTimes = new Map();
|
||||
|
||||
client.on('connect', () => {
|
||||
console.log('✅ Terhubung ke HiveMQ Cloud!');
|
||||
console.log('--------------------------------------------------');
|
||||
console.log('1. Silakan lakukan aksi di Flutter app untuk tes latensi kontrol.');
|
||||
console.log('2. Atau tekan [ENTER] di terminal ini untuk tes ping.');
|
||||
console.log('--------------------------------------------------');
|
||||
|
||||
// Subscribe ke topik sensor, perintah, dan ping
|
||||
client.subscribe([TOPIC_SENSOR, TOPIC_RELAY_WILDCARD, TOPIC_PING], (err) => {
|
||||
if (err) console.error('❌ Gagal subscribe:', err.message);
|
||||
});
|
||||
});
|
||||
|
||||
client.on('message', (topic, payload) => {
|
||||
const now = Date.now();
|
||||
const timeStr = new Date().toLocaleTimeString();
|
||||
|
||||
// 1. Logika untuk Ping Loopback (Mengukur Backend -> MQTT dan MQTT -> Backend secara presisi)
|
||||
if (topic === TOPIC_PING) {
|
||||
try {
|
||||
const data = JSON.parse(payload.toString());
|
||||
const { pingId } = data;
|
||||
const t_receive = now;
|
||||
|
||||
const t0 = pingSentTimes.get(pingId);
|
||||
if (t0) {
|
||||
const totalRTT = t_receive - t0; // Total waktu bolak-balik
|
||||
const oneWayLatency = Math.round(totalRTT / 2); // Estimasi satu arah
|
||||
|
||||
console.log(`\n[${timeStr}] ⚡ HASIL PENGUKURAN LATENSI:`);
|
||||
console.log(` 🔸 1. Backend ➔ MQTT Broker (Kirim Perintah) : ${oneWayLatency} ms`);
|
||||
console.log(` 🔸 2. MQTT Broker ➔ Backend (Terima Data ESP) : ${oneWayLatency} ms`);
|
||||
console.log(` ➔ Total Round-Trip Time (RTT) : ${totalRTT} ms`);
|
||||
console.log(`--------------------------------------------------`);
|
||||
|
||||
// Bersihkan memori cache
|
||||
pingSentTimes.delete(pingId);
|
||||
pingPubAckTimes.delete(pingId);
|
||||
}
|
||||
} catch (e) {}
|
||||
return;
|
||||
}
|
||||
|
||||
// 2. Jika menerima pesan perintah relay
|
||||
if (topic.startsWith('cmd/relay/')) {
|
||||
const deviceId = topic.split('/')[2];
|
||||
const state = payload.toString().trim();
|
||||
console.log(`[${timeStr}] 📤 Perintah [${state}] terkirim ke device: ${deviceId}`);
|
||||
pendingCommands.set(deviceId, {
|
||||
targetState: state === 'ON',
|
||||
sentAt: now
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
// 3. Jika menerima data sensor dari ESP32
|
||||
if (topic === TOPIC_SENSOR) {
|
||||
try {
|
||||
const data = JSON.parse(payload.toString());
|
||||
const deviceId = data.device_id || 'esp32-01';
|
||||
const relayState = data.relay_state ?? data.relay ?? false;
|
||||
|
||||
console.log(`[${timeStr}] 📥 Data sensor masuk dari ${deviceId} (Suhu: ${data.temp || data.temperature}°C | Hum: ${data.hum || data.humidity}% | Relay: ${relayState})`);
|
||||
|
||||
// Cek apakah ada perintah relay yang sedang ditunggu responnya
|
||||
const pending = pendingCommands.get(deviceId);
|
||||
if (pending && pending.targetState === relayState) {
|
||||
const latency = now - pending.sentAt;
|
||||
console.log(`🔥 [Control RTT] Perintah dikirim hingga status terkonfirmasi: ${latency} ms`);
|
||||
console.log(` └─ Estimasi Latensi Kirim Perintah (Backend ➔ ESP32): ${Math.round(latency / 2)} ms`);
|
||||
pendingCommands.delete(deviceId);
|
||||
}
|
||||
} catch (e) {
|
||||
console.log(`[${timeStr}] 📥 Data sensor masuk (Bukan JSON):`, payload.toString());
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
client.on('error', (err) => {
|
||||
console.error('❌ MQTT Error:', err.message);
|
||||
});
|
||||
|
||||
// Setup agar saat menekan Enter di terminal akan memicu publish ping ke HiveMQ
|
||||
process.stdin.on('data', () => {
|
||||
if (client.connected) {
|
||||
const pingId = Math.random().toString(36).substring(7);
|
||||
const t0 = Date.now();
|
||||
pingSentTimes.set(pingId, t0);
|
||||
|
||||
// Gunakan QoS 1 agar mendapatkan callback PUBACK dari broker saat pesan sukses diterima broker
|
||||
client.publish(TOPIC_PING, JSON.stringify({ pingId }), { qos: 1 }, (err) => {
|
||||
if (!err) {
|
||||
const t_puback = Date.now();
|
||||
pingPubAckTimes.set(pingId, t_puback);
|
||||
} else {
|
||||
console.error('❌ Gagal publish ping:', err.message);
|
||||
}
|
||||
});
|
||||
}
|
||||
});
|
||||
Loading…
Reference in New Issue