Spaces:
Sleeping
Sleeping
| 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(); |