feat: ApsGo Railway Worker - 24/7 Background Automation
Complete Node.js worker for IoT automation: Features: - Waktu Mode: Time-based scheduling with cron - Sensor Mode: Threshold-based automation - BullMQ + Redis: Queue management (prevent race conditions) - Firebase Admin: Realtime Database integration - Auto History Logging: Every 10 minutes - Health Monitoring: Auto-restart on failure - Graceful Shutdown: Clean resource cleanup Architecture: - NO HTTP server (background worker only) - Listen to Firebase Realtime DB changes - Process jobs from Redis queue - Execute scheduled tasks via cron Tech Stack: - Node.js 18+ - Firebase Admin SDK - BullMQ (job queue) - Redis (in-memory DB) - IoRedis client Files: - worker.js: Main worker implementation (500+ lines) - package.json: Dependencies and scripts - railway.json: Railway deployment config (NO healthcheck - worker not HTTP) - .env.example: Environment variables template - README.md: Complete documentation - .gitignore: Node modules, env files, logs Ready for Railway deployment! Connect to Redis and set Firebase environment variables.
This commit is contained in:
parent
8df97e29ec
commit
a9d894b2fd
|
|
@ -0,0 +1,13 @@
|
|||
# Firebase Configuration
|
||||
FIREBASE_PROJECT_ID=your-project-id
|
||||
FIREBASE_CLIENT_EMAIL=firebase-adminsdk-xxxxx@your-project-id.iam.gserviceaccount.com
|
||||
FIREBASE_PRIVATE_KEY="-----BEGIN PRIVATE KEY-----\nYourPrivateKeyHere\n-----END PRIVATE KEY-----\n"
|
||||
FIREBASE_DATABASE_URL=https://your-project-id-default-rtdb.firebaseio.com
|
||||
|
||||
# Redis Configuration (Railway will auto-provide this)
|
||||
REDIS_HOST=redis.railway.internal
|
||||
REDIS_PORT=6379
|
||||
REDIS_PASSWORD=
|
||||
|
||||
# Or use REDIS_URL (Railway format)
|
||||
# REDIS_URL=redis://default:password@redis.railway.internal:6379
|
||||
|
|
@ -0,0 +1,27 @@
|
|||
# Dependencies
|
||||
node_modules/
|
||||
package-lock.json
|
||||
yarn.lock
|
||||
|
||||
# Environment variables
|
||||
.env
|
||||
.env.local
|
||||
|
||||
# Logs
|
||||
logs/
|
||||
*.log
|
||||
npm-debug.log*
|
||||
|
||||
# OS files
|
||||
.DS_Store
|
||||
Thumbs.db
|
||||
|
||||
# IDE
|
||||
.vscode/
|
||||
.idea/
|
||||
*.swp
|
||||
*.swo
|
||||
|
||||
# Firebase service account
|
||||
serviceAccount.json
|
||||
firebase-key.json
|
||||
149
README.md
149
README.md
|
|
@ -1 +1,148 @@
|
|||
# myreppril
|
||||
# ApsGo Railway Worker
|
||||
|
||||
Background worker service untuk sistem otomasi IoT ApsGo. Service ini berjalan 24/7 di cloud untuk menjalankan penjadwalan dan automation bahkan ketika aplikasi mobile ditutup atau handphone pengguna mati.
|
||||
|
||||
## Features
|
||||
|
||||
- ✅ **Waktu Mode**: Penjadwalan berdasarkan waktu (cron-based)
|
||||
- ✅ **Sensor Mode**: Otomasi berdasarkan threshold kelembapan tanah
|
||||
- ✅ **Auto History Logging**: Record data sensor setiap 10 menit
|
||||
- ✅ **Redis Queue**: Prevent race conditions dan manage concurrent tasks
|
||||
- ✅ **Graceful Shutdown**: Clean shutdown dengan safety turn-off semua aktuator
|
||||
- ✅ **Health Monitoring**: Auto health check setiap 5 menit
|
||||
- ✅ **Auto Cleanup**: Hapus history lama otomatis (retain 30 hari)
|
||||
|
||||
## Tech Stack
|
||||
|
||||
- **Node.js**: Runtime environment
|
||||
- **Firebase Admin SDK**: Realtime Database integration
|
||||
- **BullMQ**: Robust job queue dengan Redis
|
||||
- **Redis**: In-memory database untuk queue dan caching
|
||||
- **Cron**: Scheduled tasks
|
||||
|
||||
## Setup Local Development
|
||||
|
||||
1. Install dependencies:
|
||||
```bash
|
||||
npm install
|
||||
```
|
||||
|
||||
2. Copy `.env.example` ke `.env` dan isi dengan credentials Firebase Anda:
|
||||
```bash
|
||||
cp .env.example .env
|
||||
```
|
||||
|
||||
3. Setup Redis lokal (gunakan Docker):
|
||||
```bash
|
||||
docker run -d -p 6379:6379 redis:latest
|
||||
```
|
||||
|
||||
4. Run worker:
|
||||
```bash
|
||||
npm run dev # Development mode dengan nodemon
|
||||
# atau
|
||||
npm start # Production mode
|
||||
```
|
||||
|
||||
## Deploy to Railway
|
||||
|
||||
Lihat file `DEPLOYMENT_GUIDE.md` untuk step-by-step deployment ke Railway.
|
||||
|
||||
## Environment Variables
|
||||
|
||||
| Variable | Description | Required |
|
||||
|----------|-------------|----------|
|
||||
| `FIREBASE_PROJECT_ID` | Firebase project ID | ✅ |
|
||||
| `FIREBASE_CLIENT_EMAIL` | Firebase service account email | ✅ |
|
||||
| `FIREBASE_PRIVATE_KEY` | Firebase service account private key | ✅ |
|
||||
| `FIREBASE_DATABASE_URL` | Firebase Realtime Database URL | ✅ |
|
||||
| `REDIS_HOST` | Redis hostname | ✅ |
|
||||
| `REDIS_PORT` | Redis port (default: 6379) | ❌ |
|
||||
| `REDIS_PASSWORD` | Redis password (if required) | ❌ |
|
||||
|
||||
## Architecture
|
||||
|
||||
```
|
||||
Flutter App (Mobile)
|
||||
↕
|
||||
Firebase Realtime DB ← ESP32/Hardware
|
||||
↕
|
||||
Railway Worker (This service)
|
||||
↕
|
||||
Redis Queue
|
||||
```
|
||||
|
||||
## How It Works
|
||||
|
||||
### Waktu Mode
|
||||
- Worker check Firebase `/kontrol` setiap 30 detik
|
||||
- Jika `waktu_1` atau `waktu_2` match dengan waktu sekarang, add job ke queue
|
||||
- Job akan diprocess oleh worker untuk nyalakan pompa dan valve
|
||||
- Setelah durasi selesai, otomatis matikan
|
||||
|
||||
### Sensor Mode
|
||||
- Worker listen ke Firebase `/data` secara realtime
|
||||
- Jika `soil_X` < `batas_bawah`, trigger watering untuk pot tersebut
|
||||
- Ada cooldown 2 menit per pot untuk prevent over-watering
|
||||
- Support 2 mode: `fixed` (durasi tetap) dan `smart` (sampai mencapai batas_atas)
|
||||
|
||||
### Safety Features
|
||||
- Concurrency: 1 (hanya 1 job diprocess pada satu waktu)
|
||||
- Debouncing: Minimum 2 menit antar penyiraman per pot
|
||||
- Error handling: Jika error, otomatis turn OFF semua aktuator
|
||||
- Graceful shutdown: Clean up resources saat restart/shutdown
|
||||
|
||||
## Monitoring
|
||||
|
||||
Worker akan log semua aktivitas ke console:
|
||||
- ✅ Success operations
|
||||
- ❌ Errors dengan details
|
||||
- 💧 Watering jobs progress
|
||||
- 📊 History logging
|
||||
- 💚 Health check status
|
||||
|
||||
Di Railway dashboard, Anda bisa:
|
||||
- View logs realtime
|
||||
- Monitor CPU/Memory usage
|
||||
- Setup alerts untuk failures
|
||||
|
||||
## Maintenance
|
||||
|
||||
### Manual Queue Management
|
||||
|
||||
Untuk clear queue (jika ada masalah):
|
||||
```javascript
|
||||
const { Queue } = require('bullmq');
|
||||
const Redis = require('ioredis');
|
||||
|
||||
const redis = new Redis(process.env.REDIS_URL);
|
||||
const queue = new Queue('watering', { connection: redis });
|
||||
|
||||
// Clear all jobs
|
||||
await queue.obliterate();
|
||||
```
|
||||
|
||||
### Database Cleanup
|
||||
|
||||
History otomatis di-cleanup setiap hari jam 2 pagi, hanya retain 30 hari terakhir.
|
||||
|
||||
## Troubleshooting
|
||||
|
||||
### Worker tidak berjalan
|
||||
1. Check environment variables
|
||||
2. Check Firebase credentials
|
||||
3. Check Redis connection
|
||||
|
||||
### Job tidak diprocess
|
||||
1. Check queue status di logs
|
||||
2. Verify Firebase rules mengizinkan admin access
|
||||
3. Check concurrency setting
|
||||
|
||||
### Memory leak
|
||||
- Worker menggunakan BullMQ yang sudah optimize untuk long-running process
|
||||
- Auto cleanup completed jobs (retain last 100)
|
||||
- Auto cleanup failed jobs (retain last 50)
|
||||
|
||||
## License
|
||||
|
||||
MIT
|
||||
|
|
|
|||
|
|
@ -0,0 +1,27 @@
|
|||
{
|
||||
"name": "apsgo-railway-worker",
|
||||
"version": "1.0.0",
|
||||
"description": "Background worker for ApsGo IoT automation system",
|
||||
"main": "worker.js",
|
||||
"scripts": {
|
||||
"start": "node worker.js",
|
||||
"dev": "nodemon worker.js",
|
||||
"test": "echo \"No tests yet\" && exit 0"
|
||||
},
|
||||
"keywords": ["iot", "automation", "firebase", "worker"],
|
||||
"author": "ApsGo Team",
|
||||
"license": "MIT",
|
||||
"dependencies": {
|
||||
"firebase-admin": "^12.0.0",
|
||||
"bullmq": "^5.1.0",
|
||||
"ioredis": "^5.3.2",
|
||||
"dotenv": "^16.4.5",
|
||||
"cron": "^3.1.6"
|
||||
},
|
||||
"devDependencies": {
|
||||
"nodemon": "^3.0.3"
|
||||
},
|
||||
"engines": {
|
||||
"node": ">=18.0.0"
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,13 @@
|
|||
{
|
||||
"$schema": "https://railway.app/railway.schema.json",
|
||||
"build": {
|
||||
"builder": "NIXPACKS",
|
||||
"buildCommand": "npm install"
|
||||
},
|
||||
"deploy": {
|
||||
"numReplicas": 1,
|
||||
"restartPolicyType": "ON_FAILURE",
|
||||
"restartPolicyMaxRetries": 10,
|
||||
"startCommand": "npm start"
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,498 @@
|
|||
/**
|
||||
* ApsGo Railway Worker
|
||||
* Background service untuk automation scheduling 24/7
|
||||
* Features:
|
||||
* - Waktu Mode: Scheduled watering by time
|
||||
* - Sensor Mode: Automatic watering by soil moisture threshold
|
||||
* - Redis Queue: Prevent race conditions & concurrent task management
|
||||
* - Firebase Realtime DB: Sync dengan Flutter app dan ESP32
|
||||
*/
|
||||
|
||||
require('dotenv').config();
|
||||
const admin = require('firebase-admin');
|
||||
const { Queue, Worker } = require('bullmq');
|
||||
const Redis = require('ioredis');
|
||||
const cron = require('cron');
|
||||
|
||||
// ==================== CONFIGURATION ====================
|
||||
|
||||
const config = {
|
||||
redis: {
|
||||
host: process.env.REDIS_HOST || 'localhost',
|
||||
port: parseInt(process.env.REDIS_PORT) || 6379,
|
||||
password: process.env.REDIS_PASSWORD || undefined,
|
||||
maxRetriesPerRequest: null, // Required for BullMQ
|
||||
},
|
||||
firebase: {
|
||||
projectId: process.env.FIREBASE_PROJECT_ID,
|
||||
clientEmail: process.env.FIREBASE_CLIENT_EMAIL,
|
||||
privateKey: process.env.FIREBASE_PRIVATE_KEY?.replace(/\\n/g, '\n'),
|
||||
databaseURL: process.env.FIREBASE_DATABASE_URL,
|
||||
},
|
||||
worker: {
|
||||
concurrency: 1, // Process 1 job at a time (prevent race condition)
|
||||
checkInterval: 30000, // Check jadwal setiap 30 detik
|
||||
sensorDebounce: 120000, // 2 menit minimum antar penyiraman per pot
|
||||
},
|
||||
};
|
||||
|
||||
console.log('🚀 Starting ApsGo Railway Worker...');
|
||||
console.log(`📡 Firebase Project: ${config.firebase.projectId}`);
|
||||
console.log(`📦 Redis: ${config.redis.host}:${config.redis.port}`);
|
||||
|
||||
// ==================== FIREBASE INITIALIZATION ====================
|
||||
|
||||
try {
|
||||
admin.initializeApp({
|
||||
credential: admin.credential.cert({
|
||||
projectId: config.firebase.projectId,
|
||||
clientEmail: config.firebase.clientEmail,
|
||||
privateKey: config.firebase.privateKey,
|
||||
}),
|
||||
databaseURL: config.firebase.databaseURL,
|
||||
});
|
||||
console.log('✅ Firebase Admin initialized');
|
||||
} catch (error) {
|
||||
console.error('❌ Firebase initialization failed:', error.message);
|
||||
process.exit(1);
|
||||
}
|
||||
|
||||
const db = admin.database();
|
||||
|
||||
// ==================== REDIS & QUEUE SETUP ====================
|
||||
|
||||
const redis = new Redis(config.redis);
|
||||
const wateringQueue = new Queue('watering', { connection: redis });
|
||||
|
||||
redis.on('connect', () => console.log('✅ Redis connected'));
|
||||
redis.on('error', (err) => console.error('❌ Redis error:', err.message));
|
||||
|
||||
// Track last watering time untuk prevent spam
|
||||
const lastWateringTime = {};
|
||||
|
||||
// ==================== WATERING WORKER ====================
|
||||
|
||||
const wateringWorker = new Worker(
|
||||
'watering',
|
||||
async (job) => {
|
||||
const { type, potNumbers, pompaAir, pompaPupuk, duration, scheduleId } = job.data;
|
||||
|
||||
console.log(`\n💧 Processing Job: ${job.id}`);
|
||||
console.log(` Type: ${type}`);
|
||||
console.log(` Pots: [${potNumbers.join(', ')}]`);
|
||||
console.log(` Duration: ${duration}s`);
|
||||
|
||||
try {
|
||||
// Prepare aktuator updates
|
||||
const updates = {};
|
||||
if (pompaAir) updates['mosvet_1'] = true;
|
||||
if (pompaPupuk) updates['mosvet_2'] = true;
|
||||
|
||||
// Turn ON valves for selected pots
|
||||
for (const pot of potNumbers) {
|
||||
if (pot >= 1 && pot <= 5) {
|
||||
updates[`mosvet_${pot + 2}`] = true; // pot 1 → mosvet_3, etc.
|
||||
}
|
||||
}
|
||||
|
||||
// Turn ON
|
||||
console.log(' 🔛 Turning ON:', Object.keys(updates).join(', '));
|
||||
await db.ref('aktuator').update(updates);
|
||||
|
||||
// Wait for duration with progress logging
|
||||
const startTime = Date.now();
|
||||
const endTime = startTime + duration * 1000;
|
||||
|
||||
while (Date.now() < endTime) {
|
||||
const remaining = Math.ceil((endTime - Date.now()) / 1000);
|
||||
if (remaining % 10 === 0 || remaining <= 5) {
|
||||
console.log(` ⏳ ${remaining}s remaining...`);
|
||||
}
|
||||
await sleep(1000);
|
||||
}
|
||||
|
||||
// Turn OFF
|
||||
const offUpdates = {};
|
||||
for (const key in updates) {
|
||||
offUpdates[key] = false;
|
||||
}
|
||||
console.log(' 🔴 Turning OFF');
|
||||
await db.ref('aktuator').update(offUpdates);
|
||||
|
||||
// Log history
|
||||
await logHistory(type, potNumbers, duration);
|
||||
|
||||
// Update last watering time
|
||||
for (const pot of potNumbers) {
|
||||
lastWateringTime[`pot_${pot}`] = Date.now();
|
||||
}
|
||||
|
||||
console.log(` ✅ Job completed successfully`);
|
||||
return { success: true, duration, pots: potNumbers };
|
||||
} catch (error) {
|
||||
console.error(` ❌ Job failed:`, error.message);
|
||||
|
||||
// Safety: Turn OFF everything
|
||||
try {
|
||||
await db.ref('aktuator').update({
|
||||
mosvet_1: false,
|
||||
mosvet_2: false,
|
||||
mosvet_3: false,
|
||||
mosvet_4: false,
|
||||
mosvet_5: false,
|
||||
mosvet_6: false,
|
||||
mosvet_7: false,
|
||||
mosvet_8: false, // Pengaduk
|
||||
});
|
||||
console.log(' 🛡️ Safety: All aktuators turned OFF');
|
||||
} catch (safetyError) {
|
||||
console.error(' ⚠️ Safety OFF failed:', safetyError.message);
|
||||
}
|
||||
|
||||
throw error;
|
||||
}
|
||||
},
|
||||
{
|
||||
connection: redis,
|
||||
concurrency: config.worker.concurrency,
|
||||
removeOnComplete: { count: 100 }, // Keep last 100 completed jobs
|
||||
removeOnFail: { count: 50 }, // Keep last 50 failed jobs
|
||||
}
|
||||
);
|
||||
|
||||
wateringWorker.on('completed', (job) => {
|
||||
console.log(`✅ Worker completed job ${job.id}`);
|
||||
});
|
||||
|
||||
wateringWorker.on('failed', (job, err) => {
|
||||
console.error(`❌ Worker failed job ${job?.id}:`, err.message);
|
||||
});
|
||||
|
||||
// ==================== WAKTU MODE (TIME SCHEDULER) ====================
|
||||
|
||||
let lastScheduleCheck = {};
|
||||
|
||||
async function checkScheduledWatering() {
|
||||
try {
|
||||
const snapshot = await db.ref('kontrol').once('value');
|
||||
const kontrolConfig = snapshot.val();
|
||||
|
||||
if (!kontrolConfig || !kontrolConfig.waktu) {
|
||||
// Waktu mode disabled
|
||||
return;
|
||||
}
|
||||
|
||||
const now = new Date();
|
||||
const currentTime = `${now.getHours().toString().padStart(2, '0')}:${now.getMinutes().toString().padStart(2, '0')}`;
|
||||
const dateKey = `${now.getFullYear()}-${(now.getMonth() + 1).toString().padStart(2, '0')}-${now.getDate().toString().padStart(2, '0')}`;
|
||||
|
||||
// Check Jadwal 1
|
||||
if (kontrolConfig.waktu_1 && kontrolConfig.waktu_1 === currentTime) {
|
||||
const scheduleKey = `jadwal_1_${dateKey}_${currentTime}`;
|
||||
|
||||
if (!lastScheduleCheck[scheduleKey]) {
|
||||
console.log(`\n🕐 JADWAL 1 TRIGGERED: ${currentTime}`);
|
||||
|
||||
await wateringQueue.add(
|
||||
'schedule-1',
|
||||
{
|
||||
type: 'waktu_jadwal_1',
|
||||
potNumbers: [1, 2, 3, 4, 5], // All pots
|
||||
pompaAir: true,
|
||||
pompaPupuk: true,
|
||||
duration: kontrolConfig.durasi_1 || 60,
|
||||
scheduleId: scheduleKey,
|
||||
},
|
||||
{
|
||||
jobId: scheduleKey,
|
||||
removeOnComplete: true,
|
||||
}
|
||||
);
|
||||
|
||||
lastScheduleCheck[scheduleKey] = true;
|
||||
console.log(` 📌 Added to queue: ${scheduleKey}`);
|
||||
}
|
||||
}
|
||||
|
||||
// Check Jadwal 2
|
||||
if (kontrolConfig.waktu_2 && kontrolConfig.waktu_2 === currentTime) {
|
||||
const scheduleKey = `jadwal_2_${dateKey}_${currentTime}`;
|
||||
|
||||
if (!lastScheduleCheck[scheduleKey]) {
|
||||
console.log(`\n🕑 JADWAL 2 TRIGGERED: ${currentTime}`);
|
||||
|
||||
await wateringQueue.add(
|
||||
'schedule-2',
|
||||
{
|
||||
type: 'waktu_jadwal_2',
|
||||
potNumbers: [1, 2, 3, 4, 5], // All pots
|
||||
pompaAir: true,
|
||||
pompaPupuk: true,
|
||||
duration: kontrolConfig.durasi_2 || 60,
|
||||
scheduleId: scheduleKey,
|
||||
},
|
||||
{
|
||||
jobId: scheduleKey,
|
||||
removeOnComplete: true,
|
||||
}
|
||||
);
|
||||
|
||||
lastScheduleCheck[scheduleKey] = true;
|
||||
console.log(` 📌 Added to queue: ${scheduleKey}`);
|
||||
}
|
||||
}
|
||||
|
||||
// Cleanup old schedule checks (> 2 menit)
|
||||
const twoMinutesAgo = Date.now() - 120000;
|
||||
for (const key in lastScheduleCheck) {
|
||||
if (key.includes(dateKey)) continue; // Keep today's
|
||||
delete lastScheduleCheck[key];
|
||||
}
|
||||
} catch (error) {
|
||||
console.error('❌ Error checking scheduled watering:', error.message);
|
||||
}
|
||||
}
|
||||
|
||||
// Run check setiap 30 detik
|
||||
setInterval(checkScheduledWatering, config.worker.checkInterval);
|
||||
console.log(`✅ Waktu Mode scheduler started (check every ${config.worker.checkInterval / 1000}s)`);
|
||||
|
||||
// ==================== SENSOR MODE (THRESHOLD MONITORING) ====================
|
||||
|
||||
async function setupSensorMonitoring() {
|
||||
console.log('✅ Sensor Mode monitoring started');
|
||||
|
||||
db.ref('data').on('value', async (snapshot) => {
|
||||
try {
|
||||
const sensorData = snapshot.val();
|
||||
if (!sensorData) return;
|
||||
|
||||
const configSnapshot = await db.ref('kontrol').once('value');
|
||||
const kontrolConfig = configSnapshot.val();
|
||||
|
||||
if (!kontrolConfig || !kontrolConfig.otomatis) {
|
||||
// Sensor mode disabled
|
||||
return;
|
||||
}
|
||||
|
||||
const batasBawah = kontrolConfig.batas_bawah || 40;
|
||||
const batasAtas = kontrolConfig.batas_atas || 100;
|
||||
const durasiSensor = kontrolConfig.durasi_sensor || 60;
|
||||
const modeSensor = kontrolConfig.mode_sensor || 'fixed'; // 'fixed' or 'smart'
|
||||
|
||||
// Check each pot
|
||||
for (let i = 1; i <= 5; i++) {
|
||||
const soilKey = `soil_${i}`;
|
||||
const soilValue = parseInt(sensorData[soilKey]) || 0;
|
||||
|
||||
if (soilValue < batasBawah) {
|
||||
const potKey = `pot_${i}`;
|
||||
const lastTime = lastWateringTime[potKey];
|
||||
|
||||
// Debounce: minimum 2 menit antar penyiraman
|
||||
if (lastTime && Date.now() - lastTime < config.worker.sensorDebounce) {
|
||||
const remainingSeconds = Math.ceil((config.worker.sensorDebounce - (Date.now() - lastTime)) / 1000);
|
||||
console.log(`⏳ POT ${i}: Cooldown active (${remainingSeconds}s remaining)`);
|
||||
continue;
|
||||
}
|
||||
|
||||
console.log(`\n🌡️ SENSOR TRIGGERED: POT ${i}`);
|
||||
console.log(` Soil moisture: ${soilValue}% < ${batasBawah}%`);
|
||||
console.log(` Mode: ${modeSensor}, Duration: ${durasiSensor}s`);
|
||||
|
||||
const jobId = `sensor-pot-${i}-${Date.now()}`;
|
||||
await wateringQueue.add(
|
||||
`sensor-pot-${i}`,
|
||||
{
|
||||
type: 'sensor_threshold',
|
||||
potNumbers: [i],
|
||||
pompaAir: true,
|
||||
pompaPupuk: false, // No pupuk for sensor mode
|
||||
duration: durasiSensor,
|
||||
scheduleId: jobId,
|
||||
sensorData: { soilValue, batasBawah, batasAtas, mode: modeSensor },
|
||||
},
|
||||
{
|
||||
jobId,
|
||||
removeOnComplete: true,
|
||||
priority: 1, // Higher priority for sensor-triggered
|
||||
}
|
||||
);
|
||||
|
||||
console.log(` 📌 Added to queue: ${jobId}`);
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
console.error('❌ Error in sensor monitoring:', error.message);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
setupSensorMonitoring();
|
||||
|
||||
// ==================== HISTORY LOGGING ====================
|
||||
|
||||
async function logHistory(type, potNumbers, duration) {
|
||||
try {
|
||||
const now = new Date();
|
||||
const dateKey = `${now.getFullYear()}-${(now.getMonth() + 1).toString().padStart(2, '0')}-${now.getDate().toString().padStart(2, '0')}`;
|
||||
const timeKey = `${now.getHours().toString().padStart(2, '0')}:${now.getMinutes().toString().padStart(2, '0')}`;
|
||||
|
||||
// Get current sensor data
|
||||
const sensorSnapshot = await db.ref('data').once('value');
|
||||
const sensorData = sensorSnapshot.val() || {};
|
||||
|
||||
await db.ref(`history/${dateKey}/${timeKey}`).set({
|
||||
timestamp: now.getTime(),
|
||||
type: type,
|
||||
pots: potNumbers,
|
||||
duration: duration,
|
||||
...sensorData,
|
||||
});
|
||||
|
||||
console.log(` 📊 History logged: ${dateKey} ${timeKey}`);
|
||||
} catch (error) {
|
||||
console.error(' ⚠️ Failed to log history:', error.message);
|
||||
}
|
||||
}
|
||||
|
||||
// ==================== PERIODIC HISTORY LOGGING ====================
|
||||
|
||||
// Auto-log sensor data setiap 10 menit (independent from watering)
|
||||
const autoLogJob = new cron.CronJob('*/10 * * * *', async () => {
|
||||
try {
|
||||
const sensorSnapshot = await db.ref('data').once('value');
|
||||
const sensorData = sensorSnapshot.val();
|
||||
|
||||
if (sensorData) {
|
||||
const now = new Date();
|
||||
const dateKey = `${now.getFullYear()}-${(now.getMonth() + 1).toString().padStart(2, '0')}-${now.getDate().toString().padStart(2, '0')}`;
|
||||
const timeKey = `${now.getHours().toString().padStart(2, '0')}:${now.getMinutes().toString().padStart(2, '0')}`;
|
||||
|
||||
await db.ref(`history/${dateKey}/${timeKey}`).set({
|
||||
timestamp: now.getTime(),
|
||||
type: 'auto_log',
|
||||
...sensorData,
|
||||
});
|
||||
|
||||
console.log(`📊 Auto-logged sensor data: ${timeKey}`);
|
||||
}
|
||||
} catch (error) {
|
||||
console.error('❌ Auto-log failed:', error.message);
|
||||
}
|
||||
});
|
||||
|
||||
autoLogJob.start();
|
||||
console.log('✅ Auto history logging started (every 10 minutes)');
|
||||
|
||||
// ==================== CLEANUP OLD HISTORY (DAILY) ====================
|
||||
|
||||
const cleanupJob = new cron.CronJob('0 2 * * *', async () => {
|
||||
// Run daily at 2 AM
|
||||
try {
|
||||
console.log('\n🧹 Running history cleanup...');
|
||||
const daysToKeep = 30;
|
||||
const cutoffDate = new Date();
|
||||
cutoffDate.setDate(cutoffDate.getDate() - daysToKeep);
|
||||
|
||||
const historySnapshot = await db.ref('history').once('value');
|
||||
const historyData = historySnapshot.val();
|
||||
|
||||
if (historyData) {
|
||||
let deletedCount = 0;
|
||||
for (const dateKey in historyData) {
|
||||
try {
|
||||
const [year, month, day] = dateKey.split('-').map(Number);
|
||||
const date = new Date(year, month - 1, day);
|
||||
|
||||
if (date < cutoffDate) {
|
||||
await db.ref(`history/${dateKey}`).remove();
|
||||
deletedCount++;
|
||||
console.log(` 🗑️ Deleted: ${dateKey}`);
|
||||
}
|
||||
} catch (error) {
|
||||
console.error(` ⚠️ Error deleting ${dateKey}:`, error.message);
|
||||
}
|
||||
}
|
||||
console.log(`✅ Cleanup completed: ${deletedCount} dates removed`);
|
||||
}
|
||||
} catch (error) {
|
||||
console.error('❌ Cleanup failed:', error.message);
|
||||
}
|
||||
});
|
||||
|
||||
cleanupJob.start();
|
||||
console.log('✅ History cleanup scheduled (daily at 2 AM)');
|
||||
|
||||
// ==================== UTILITIES ====================
|
||||
|
||||
function sleep(ms) {
|
||||
return new Promise((resolve) => setTimeout(resolve, ms));
|
||||
}
|
||||
|
||||
// ==================== HEALTH CHECK ====================
|
||||
|
||||
async function healthCheck() {
|
||||
try {
|
||||
// Check Firebase connection
|
||||
await db.ref('.info/connected').once('value');
|
||||
|
||||
// Check Redis connection
|
||||
await redis.ping();
|
||||
|
||||
// Check queue
|
||||
const queueStatus = await wateringQueue.getJobCounts();
|
||||
|
||||
console.log('\n💚 HEALTH CHECK:');
|
||||
console.log(` Firebase: ✅ Connected`);
|
||||
console.log(` Redis: ✅ Connected`);
|
||||
console.log(` Queue: ${queueStatus.active} active, ${queueStatus.waiting} waiting`);
|
||||
} catch (error) {
|
||||
console.error('❤️🩹 HEALTH CHECK FAILED:', error.message);
|
||||
}
|
||||
}
|
||||
|
||||
// Run health check every 5 minutes
|
||||
setInterval(healthCheck, 300000);
|
||||
|
||||
// ==================== GRACEFUL SHUTDOWN ====================
|
||||
|
||||
async function shutdown() {
|
||||
console.log('\n🛑 Shutting down gracefully...');
|
||||
|
||||
try {
|
||||
await wateringWorker.close();
|
||||
console.log('✅ Worker closed');
|
||||
|
||||
await wateringQueue.close();
|
||||
console.log('✅ Queue closed');
|
||||
|
||||
await redis.quit();
|
||||
console.log('✅ Redis disconnected');
|
||||
|
||||
await admin.app().delete();
|
||||
console.log('✅ Firebase disconnected');
|
||||
|
||||
process.exit(0);
|
||||
} catch (error) {
|
||||
console.error('❌ Shutdown error:', error.message);
|
||||
process.exit(1);
|
||||
}
|
||||
}
|
||||
|
||||
process.on('SIGTERM', shutdown);
|
||||
process.on('SIGINT', shutdown);
|
||||
|
||||
// ==================== STARTUP COMPLETE ====================
|
||||
|
||||
console.log('\n✨ ApsGo Railway Worker is running!');
|
||||
console.log('📊 Features enabled:');
|
||||
console.log(' • Waktu Mode (Time-based scheduling)');
|
||||
console.log(' • Sensor Mode (Threshold-based automation)');
|
||||
console.log(' • Auto History Logging (every 10 min)');
|
||||
console.log(' • History Cleanup (daily at 2 AM)');
|
||||
console.log(' • Health Check (every 5 min)');
|
||||
console.log('\n🎯 Worker is ready to process jobs...\n');
|
||||
|
||||
// Initial health check
|
||||
setTimeout(healthCheck, 5000);
|
||||
Loading…
Reference in New Issue