Spaces:
Sleeping
Sleeping
File size: 4,136 Bytes
b14a10f | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 | 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(); |