require('dotenv').config(); const WebSocket = require('ws'); const translate = require('@vitalets/google-translate-api'); const fs = require('fs'); const path = require('path'); const os = require('os'); const winston = require('winston'); const settings = JSON.parse(fs.readFileSync(path.join(__dirname, '../settings.json'), 'utf8')); // Logger const logger = winston.createLogger({ level: 'info', format: winston.format.combine( winston.format.timestamp(), winston.format.json() ), transports: [ new winston.transports.Console(), new winston.transports.File({ filename: path.join(__dirname, '../../logs/worker.log') }) ] }); // Worker state let workerId = null; let ws = null; let maxConcurrentJobs = settings.maxConcurrentJobs; let activeJobs = 0; let reconnectAttempts = 0; const MAX_RECONNECT_ATTEMPTS = 10; const RECONNECT_DELAY = 5000; // Connect to central server function connect() { const url = `ws://${settings.centralHost}:${settings.centralPort}`; logger.info(`Connecting to central server at ${url}`); ws = new WebSocket(url); ws.on('open', () => { logger.info('Connected to central server'); reconnectAttempts = 0; }); ws.on('message', (data) => { try { const msg = JSON.parse(data); handleMessage(msg); } catch (e) { logger.error('Invalid message from central', e); } }); ws.on('close', () => { logger.warn('Disconnected from central server'); scheduleReconnect(); }); ws.on('error', (err) => { logger.error('WebSocket error', err); }); } function scheduleReconnect() { if (reconnectAttempts >= MAX_RECONNECT_ATTEMPTS) { logger.error('Max reconnect attempts reached, exiting'); process.exit(1); } reconnectAttempts++; const delay = RECONNECT_DELAY * reconnectAttempts; logger.info(`Reconnecting in ${delay}ms (attempt ${reconnectAttempts}/${MAX_RECONNECT_ATTEMPTS})`); setTimeout(connect, delay); } function handleMessage(msg) { switch (msg.type) { case 'welcome': workerId = msg.workerId; maxConcurrentJobs = msg.maxConcurrentJobs || maxConcurrentJobs; logger.info(`Worker registered with ID: ${workerId}, max jobs: ${maxConcurrentJobs}`); startHeartbeat(); break; case 'translate': handleTranslateJob(msg); break; case 'ping': // Keep alive ws.send(JSON.stringify({ type: 'pong' })); break; } } async function handleTranslateJob(msg) { const { jobId, text, source, target } = msg; if (activeJobs >= maxConcurrentJobs) { logger.warn(`Job ${jobId} rejected: at capacity (${activeJobs}/${maxConcurrentJobs})`); sendResult(jobId, { error: 'Worker at capacity' }); return; } activeJobs++; logger.info(`Starting job ${jobId} (${activeJobs}/${maxConcurrentJobs})`); try { const result = await translate(text, { from: source, to: target }); sendResult(jobId, { translatedText: result.text }); } catch (error) { logger.error(`Job ${jobId} failed`, error); sendResult(jobId, { error: error.message }); } finally { activeJobs--; } } function sendResult(jobId, result) { if (ws && ws.readyState === WebSocket.OPEN) { ws.send(JSON.stringify({ type: 'result', jobId, result })); } } // Heartbeat with metrics function startHeartbeat() { setInterval(() => { if (ws && ws.readyState === WebSocket.OPEN) { const cpuUsage = process.cpuUsage(); const memUsage = process.memoryUsage(); ws.send(JSON.stringify({ type: 'heartbeat', cpu: (cpuUsage.user + cpuUsage.system) / 10000, // Convert to percentage-ish ram: memUsage.heapUsed, jobsActive: activeJobs })); } }, 5000); } // Handle graceful shutdown process.on('SIGTERM', () => { logger.info('SIGTERM received, shutting down gracefully'); if (ws) ws.close(); process.exit(0); }); process.on('SIGINT', () => { logger.info('SIGINT received, shutting down gracefully'); if (ws) ws.close(); process.exit(0); }); // Start logger.info('Starting Translator Worker'); connect();