NathMen12 commited on
Commit
af47e71
·
verified ·
1 Parent(s): 78fdaed

Update central/src/index.js

Browse files
Files changed (1) hide show
  1. central/src/index.js +147 -187
central/src/index.js CHANGED
@@ -11,13 +11,9 @@ const PQueue = require('p-queue').default;
11
  const os = require('os');
12
  const winston = require('winston');
13
 
14
- // Charger settings.json (chemin relatif à ce fichier)
15
- const settingsPath = path.join(__dirname, '../settings.json');
16
- const settings = JSON.parse(fs.readFileSync(settingsPath, 'utf8'));
17
 
18
- // ──────────────────────────────────────────────
19
- // Logger
20
- // ──────────────────────────────────────────────
21
  const logger = winston.createLogger({
22
  level: 'info',
23
  format: winston.format.combine(
@@ -30,65 +26,31 @@ const logger = winston.createLogger({
30
  ]
31
  });
32
 
33
- // ──────────────────────────────────────────────
34
- // Métriques & État
35
- // ──────────────────────────────────────────────
36
  const metrics = {
37
- requests: [],
38
  cpuHistory: [],
39
  ramHistory: [],
40
- workers: new Map()
41
  };
42
 
 
43
  const localQueue = new PQueue({ concurrency: settings.maxLocalJobs });
44
- const workerQueue = new PQueue({ concurrency: 10 });
45
- const pendingJobs = new Map();
46
-
47
- // ──────────────────────────────────────────────
48
- // Helpers CPU / RAM (Machine complète)
49
- // ──────────────────────────────────────────────
50
- function getSystemCpuPercent() {
51
- // Utilise load average 1min normalisé par nb cœurs (non-bloquant, pas besoin de double mesure)
52
- const cpus = os.cpus();
53
- const load1 = os.loadavg()[0];
54
- return Math.min(100, (load1 / cpus.length) * 100);
55
- }
56
-
57
- function getSystemRamUsed() {
58
- return os.totalmem() - os.freemem();
59
- }
60
-
61
- // ──────────────────────────────────────────────
62
- // Collecte métriques (1s)
63
- // ──────────────────────────────────────────────
64
- setInterval(() => {
65
- const now = Date.now();
66
- const cpu = getSystemCpuPercent();
67
- const ramUsed = getSystemRamUsed();
68
 
69
- metrics.cpuHistory.push({ timestamp: now, cpu: Math.round(cpu * 10) / 10 });
70
- metrics.ramHistory.push({ timestamp: now, used: ramUsed, total: os.totalmem() });
71
-
72
- const cutoff = now - settings.metricsWindowSeconds * 1000;
73
- metrics.cpuHistory = metrics.cpuHistory.filter(m => m.timestamp > cutoff);
74
- metrics.ramHistory = metrics.ramHistory.filter(m => m.timestamp > cutoff);
75
- metrics.requests = metrics.requests.filter(m => m.timestamp > cutoff);
76
- }, 1000);
77
-
78
- // ──────────────────────────────────────────────
79
- // Express + WS
80
- // ──────────────────────────────────────────────
81
  const app = express();
82
  const server = http.createServer(app);
83
- const wss = new WebSocketServer({ server });
84
 
85
- const ADMIN_CODE = process.env.ADMIN_ACCESS_CODE;
86
- if (!ADMIN_CODE) logger.warn('⚠️ ADMIN_ACCESS_CODE non défini dans .env');
87
 
 
88
  app.use(cors());
89
  app.use(express.json({ limit: '10mb' }));
90
  app.use(express.static(path.join(__dirname, '../../public')));
91
 
 
92
  const limiter = rateLimit({
93
  windowMs: settings.rateLimit.windowMs,
94
  max: settings.rateLimit.maxRequestsPerMinutePerIP,
@@ -99,31 +61,52 @@ const limiter = rateLimit({
99
  });
100
  app.use('/translate', limiter);
101
 
 
102
  const adminAuth = (req, res, next) => {
 
103
  const providedCode = req.headers['x-admin-code'] || req.query.admin_code;
104
- if (!ADMIN_CODE || providedCode !== ADMIN_CODE) {
105
  return res.status(401).json({ error: 'Code admin invalide' });
106
  }
107
  next();
108
  };
109
 
110
- // ──────────────────────────────────────────────
111
- // Traduction
112
- // ──────────────────────────────────────────────
113
- const translate = require('@vitalets/google-translate-api').default;
114
 
115
- async function translateText(text, source, target) {
116
- const res = await translate(text, { from: source, to: target });
117
- return res.text;
118
- }
 
 
 
 
 
 
 
 
 
 
 
 
119
 
120
- // ──────────────────────────────────────────────
121
- // WebSocket Workers
122
- // ──────────────────────────────────────────────
 
 
 
 
 
 
 
 
 
123
  wss.on('connection', (ws, req) => {
124
  const workerId = uuidv4();
125
- logger.info(`Worker connecté: ${workerId}`);
126
-
127
  const worker = {
128
  id: workerId,
129
  ws,
@@ -131,144 +114,107 @@ wss.on('connection', (ws, req) => {
131
  ram: 0,
132
  lastHeartbeat: Date.now(),
133
  jobsActive: 0,
134
- maxJobs: settings.workerMaxJobs || 2
135
  };
136
  metrics.workers.set(workerId, worker);
137
-
138
- ws.send(JSON.stringify({
139
- type: 'welcome',
140
- workerId,
141
- maxConcurrentJobs: worker.maxJobs
142
- }));
143
-
144
  ws.on('message', (data) => {
145
  try {
146
  const msg = JSON.parse(data);
147
  handleWorkerMessage(workerId, msg);
148
  } catch (e) {
149
- logger.error('Message worker invalide', { workerId, error: e.message });
150
  }
151
  });
152
-
153
- const cleanup = () => {
154
- // Réveiller les jobs en attente pour qu'ils basculent en local
155
- for (const [jobId, pending] of pendingJobs) {
156
- if (pending.workerId === workerId) {
157
- clearTimeout(pending.timeout);
158
- pending.reject(new Error('Worker déconnecté'));
159
- pendingJobs.delete(jobId);
160
- }
161
- }
162
  metrics.workers.delete(workerId);
163
- logger.info(`Worker déconnecté: ${workerId}`);
164
- };
165
-
166
- ws.on('close', cleanup);
167
  ws.on('error', (err) => {
168
- logger.error(`Erreur worker ${workerId}`, err);
169
- cleanup();
170
  });
 
 
 
 
 
 
 
171
  });
172
 
173
  function handleWorkerMessage(workerId, msg) {
174
  const worker = metrics.workers.get(workerId);
175
  if (!worker) return;
176
-
177
  switch (msg.type) {
178
  case 'heartbeat':
179
  worker.lastHeartbeat = Date.now();
180
- worker.cpu = msg.cpu ?? 0;
181
- worker.ram = msg.ram ?? 0;
182
- worker.jobsActive = msg.jobsActive ?? 0;
183
  break;
184
-
185
  case 'result':
 
186
  if (msg.jobId) {
 
187
  const pending = pendingJobs.get(msg.jobId);
188
  if (pending) {
189
- clearTimeout(pending.timeout);
190
  pending.resolve(msg.result);
191
- pendingJobs.delete(jobId);
192
  }
193
- worker.jobsActive = Math.max(0, worker.jobsActive - 1);
194
  }
 
195
  break;
196
-
197
  case 'log':
198
- logger.info(`[Worker ${workerId.substring(0,8)}] ${msg.message}`);
199
  break;
200
  }
201
  }
202
 
203
- // ──────────────────────────────────────────────
204
- // Dispatch vers worker
205
- // ──────────────────────────────────────────────
206
- function findAvailableWorker() {
207
- for (const [, worker] of metrics.workers) {
208
- if (worker.ws.readyState === 1 && worker.jobsActive < worker.maxJobs) {
209
- return worker;
210
- }
211
- }
212
- return null;
213
- }
214
-
215
- function dispatchToWorker(worker, job) {
216
- return new Promise((resolve, reject) => {
217
- const jobId = uuidv4();
218
- const timeout = setTimeout(() => {
219
- pendingJobs.delete(jobId);
220
- worker.jobsActive = Math.max(0, worker.jobsActive - 1);
221
- reject(new Error('Worker timeout (30s)'));
222
- }, 30000);
223
 
224
- pendingJobs.set(jobId, { resolve, reject, timeout, workerId: worker.id });
225
 
226
- worker.ws.send(JSON.stringify({
227
- type: 'translate',
228
- jobId,
229
- ...job
230
- }));
231
- worker.jobsActive++;
232
- });
233
- }
234
-
235
- // ──────────────────────────────────────────────
236
- // Routes API
237
- // ──────────────────────────────────────────────
238
  app.get('/api/health', (req, res) => {
239
  res.json({ status: 'ok', uptime: process.uptime() });
240
  });
241
 
 
242
  app.post('/translate', async (req, res) => {
243
  const startTime = Date.now();
244
  const { text, source = 'auto', target = 'fr' } = req.body;
245
  const ip = req.ip;
246
-
247
  if (!text || typeof text !== 'string') {
248
  return res.status(400).json({ error: 'Texte requis' });
249
  }
 
250
  if (text.length > 5000) {
251
  return res.status(400).json({ error: 'Texte trop long (max 5000 caractères)' });
252
  }
253
-
254
  try {
255
- let translatedText;
256
- const worker = findAvailableWorker();
257
-
258
- if (worker) {
259
- try {
260
- const result = await dispatchToWorker(worker, { text, source, target });
261
- translatedText = result.translatedText ?? result;
262
- } catch (workerErr) {
263
- logger.warn(`Worker échoué, fallback local: ${workerErr.message}`);
264
- translatedText = await localQueue.add(() => translateText(text, source, target));
265
- }
266
  } else {
267
- translatedText = await localQueue.add(() => translateText(text, source, target));
 
268
  }
269
-
270
  const duration = Date.now() - startTime;
271
-
 
 
 
 
272
  metrics.requests.push({
273
  timestamp: Date.now(),
274
  ip,
@@ -278,18 +224,46 @@ app.post('/translate', async (req, res) => {
278
  inputText: text.substring(0, 100),
279
  outputText: translatedText.substring(0, 100)
280
  });
281
-
282
  res.json({ translatedText, source, target, duration });
283
  } catch (error) {
284
- logger.error('Erreur traduction', { error: error.message, stack: error.stack });
285
  res.status(500).json({ error: 'Erreur de traduction' });
286
  }
287
  });
288
 
289
- // ──────────────────────────────────────────────
290
- // Admin API FORMAT COMPATIBLE FRONTEND
291
- // ──────────────────────────────────────────────
292
- app.get('/admin', (req, res) => {
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
293
  res.sendFile(path.join(__dirname, '../../public/admin.html'));
294
  });
295
 
@@ -305,24 +279,18 @@ app.get('/api/admin/metrics', adminAuth, (req, res) => {
305
  lastHeartbeat: w.lastHeartbeat
306
  });
307
  }
308
-
309
- // Point "maintenant" pour les graphiques (évite le trou de 1s)
310
- const currentCpu = getSystemCpuPercent();
311
- const currentRam = getSystemRamUsed();
312
-
313
  res.json({
314
  workers,
315
  cpuHistory: metrics.cpuHistory,
316
  ramHistory: metrics.ramHistory,
317
- requests: metrics.requests.slice(-200),
318
- requestRate: calculateRequestRate(),
319
- // Requis par le frontend pour courbes fluides
320
- currentCpu: Math.round(currentCpu * 10) / 10,
321
- currentRam: currentRam
322
  });
323
  });
324
 
325
  app.get('/api/admin/logs', adminAuth, (req, res) => {
 
326
  const logs = metrics.requests.slice(-100).map(r => ({
327
  timestamp: r.timestamp,
328
  ip: r.ip,
@@ -337,33 +305,25 @@ app.get('/api/admin/logs', adminAuth, (req, res) => {
337
 
338
  function calculateRequestRate() {
339
  const now = Date.now();
340
- return metrics.requests.filter(r => r.timestamp > now - 120000).length;
 
 
341
  }
342
 
343
- // ──────────────────────────────────────────────
344
- // Pages statiques
345
- // ──────────────────────────────────────────────
346
- app.get('/', (req, res) => res.sendFile(path.join(__dirname, '../../public/index.html')));
347
- app.get('/wiki', (req, res) => res.sendFile(path.join(__dirname, '../../public/wiki.html')));
348
 
349
- // ──────────────────────────────────────────────
350
- // Nettoyage workers inactifs (5s)
351
- // ──────────────────────────────────────────────
352
- setInterval(() => {
353
- const now = Date.now();
354
- for (const [id, worker] of metrics.workers) {
355
- if (now - worker.lastHeartbeat > 15000) {
356
- worker.ws.close(); // Déclenchera l'event 'close' et cleanup()
357
- }
358
- }
359
- }, 5000);
360
 
361
- // ──────────────────────────────────────────────
362
- // Start
363
- // ──────────────────────────────────────────────
364
  const PORT = process.env.PORT || settings.port;
365
  server.listen(PORT, '0.0.0.0', () => {
366
- logger.info(`🚀 Serveur central sur port ${PORT}`);
367
- logger.info(`📊 Admin: http://localhost:${PORT}/admin`);
368
- logger.info(`🔧 API: http://localhost:${PORT}/translate`);
369
  });
 
11
  const os = require('os');
12
  const winston = require('winston');
13
 
14
+ const settings = JSON.parse(fs.readFileSync(path.join(__dirname, '../settings.json'), 'utf8'));
 
 
15
 
16
+ // Logger configuration
 
 
17
  const logger = winston.createLogger({
18
  level: 'info',
19
  format: winston.format.combine(
 
26
  ]
27
  });
28
 
29
+ // Metrics storage
 
 
30
  const metrics = {
31
+ requests: [], // { timestamp, ip, duration, sourceLang, targetLang, inputText, outputText }
32
  cpuHistory: [],
33
  ramHistory: [],
34
+ workers: new Map() // workerId -> { id, ws, cpu, ram, lastHeartbeat, jobsActive }
35
  };
36
 
37
+ // Job queue (local processing + worker dispatch)
38
  const localQueue = new PQueue({ concurrency: settings.maxLocalJobs });
39
+ const workerQueue = new PQueue({ concurrency: 10 }); // dispatch to workers
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
40
 
41
+ // Express app
 
 
 
 
 
 
 
 
 
 
 
42
  const app = express();
43
  const server = http.createServer(app);
 
44
 
45
+ // WebSocket server for workers
46
+ const wss = new WebSocketServer({ server });
47
 
48
+ // Middleware
49
  app.use(cors());
50
  app.use(express.json({ limit: '10mb' }));
51
  app.use(express.static(path.join(__dirname, '../../public')));
52
 
53
+ // Rate limiting
54
  const limiter = rateLimit({
55
  windowMs: settings.rateLimit.windowMs,
56
  max: settings.rateLimit.maxRequestsPerMinutePerIP,
 
61
  });
62
  app.use('/translate', limiter);
63
 
64
+ // Admin authentication middleware
65
  const adminAuth = (req, res, next) => {
66
+ const adminCode = process.env.ADMIN_ACCESS_CODE;
67
  const providedCode = req.headers['x-admin-code'] || req.query.admin_code;
68
+ if (!adminCode || providedCode !== adminCode) {
69
  return res.status(401).json({ error: 'Code admin invalide' });
70
  }
71
  next();
72
  };
73
 
74
+ // Load translation function
75
+ const translate = require('@vitalets/google-translate-api');
 
 
76
 
77
+ // Metrics collection
78
+ setInterval(() => {
79
+ const cpuUsage = process.cpuUsage();
80
+ const memUsage = process.memoryUsage();
81
+ const totalMem = os.totalmem();
82
+ const freeMem = os.freemem();
83
+
84
+ metrics.cpuHistory.push({ timestamp: Date.now(), cpu: cpuUsage.user + cpuUsage.system });
85
+ metrics.ramHistory.push({ timestamp: Date.now(), used: totalMem - freeMem, total: totalMem });
86
+
87
+ // Keep only last metricsWindowSeconds
88
+ const cutoff = Date.now() - settings.metricsWindowSeconds * 1000;
89
+ metrics.cpuHistory = metrics.cpuHistory.filter(m => m.timestamp > cutoff);
90
+ metrics.ramHistory = metrics.ramHistory.filter(m => m.timestamp > cutoff);
91
+ metrics.requests = metrics.requests.filter(m => m.timestamp > cutoff);
92
+ }, 1000);
93
 
94
+ // Clean up inactive workers
95
+ setInterval(() => {
96
+ const now = Date.now();
97
+ for (const [id, worker] of metrics.workers) {
98
+ if (now - worker.lastHeartbeat > 15000) {
99
+ logger.info(`Worker ${id} disconnected (timeout)`);
100
+ metrics.workers.delete(id);
101
+ }
102
+ }
103
+ }, 5000);
104
+
105
+ // WebSocket handling for workers
106
  wss.on('connection', (ws, req) => {
107
  const workerId = uuidv4();
108
+ logger.info(`Worker connected: ${workerId}`);
109
+
110
  const worker = {
111
  id: workerId,
112
  ws,
 
114
  ram: 0,
115
  lastHeartbeat: Date.now(),
116
  jobsActive: 0,
117
+ maxJobs: 2
118
  };
119
  metrics.workers.set(workerId, worker);
120
+
 
 
 
 
 
 
121
  ws.on('message', (data) => {
122
  try {
123
  const msg = JSON.parse(data);
124
  handleWorkerMessage(workerId, msg);
125
  } catch (e) {
126
+ logger.error('Invalid worker message', e);
127
  }
128
  });
129
+
130
+ ws.on('close', () => {
131
+ logger.info(`Worker disconnected: ${workerId}`);
 
 
 
 
 
 
 
132
  metrics.workers.delete(workerId);
133
+ });
134
+
 
 
135
  ws.on('error', (err) => {
136
+ logger.error(`Worker ${workerId} error`, err);
 
137
  });
138
+
139
+ // Send welcome message with config
140
+ ws.send(JSON.stringify({
141
+ type: 'welcome',
142
+ workerId,
143
+ maxConcurrentJobs: 2
144
+ }));
145
  });
146
 
147
  function handleWorkerMessage(workerId, msg) {
148
  const worker = metrics.workers.get(workerId);
149
  if (!worker) return;
150
+
151
  switch (msg.type) {
152
  case 'heartbeat':
153
  worker.lastHeartbeat = Date.now();
154
+ worker.cpu = msg.cpu || 0;
155
+ worker.ram = msg.ram || 0;
156
+ worker.jobsActive = msg.jobsActive || 0;
157
  break;
 
158
  case 'result':
159
+ // Job completed by worker
160
  if (msg.jobId) {
161
+ // Find and resolve the waiting promise
162
  const pending = pendingJobs.get(msg.jobId);
163
  if (pending) {
 
164
  pending.resolve(msg.result);
165
+ pendingJobs.delete(msg.jobId);
166
  }
 
167
  }
168
+ worker.jobsActive = Math.max(0, worker.jobsActive - 1);
169
  break;
 
170
  case 'log':
171
+ logger.info(`Worker ${workerId} log: ${msg.message}`);
172
  break;
173
  }
174
  }
175
 
176
+ // Pending jobs waiting for worker results
177
+ const pendingJobs = new Map();
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
178
 
179
+ // API Routes
180
 
181
+ // Health check
 
 
 
 
 
 
 
 
 
 
 
182
  app.get('/api/health', (req, res) => {
183
  res.json({ status: 'ok', uptime: process.uptime() });
184
  });
185
 
186
+ // Translate endpoint
187
  app.post('/translate', async (req, res) => {
188
  const startTime = Date.now();
189
  const { text, source = 'auto', target = 'fr' } = req.body;
190
  const ip = req.ip;
191
+
192
  if (!text || typeof text !== 'string') {
193
  return res.status(400).json({ error: 'Texte requis' });
194
  }
195
+
196
  if (text.length > 5000) {
197
  return res.status(400).json({ error: 'Texte trop long (max 5000 caractères)' });
198
  }
199
+
200
  try {
201
+ // Try to dispatch to available worker first
202
+ let result;
203
+ const availableWorker = findAvailableWorker();
204
+
205
+ if (availableWorker) {
206
+ result = await dispatchToWorker(availableWorker, { text, source, target });
 
 
 
 
 
207
  } else {
208
+ // Process locally
209
+ result = await localQueue.add(() => translateText(text, source, target));
210
  }
211
+
212
  const duration = Date.now() - startTime;
213
+
214
+ // Extract translated text from worker result object if needed
215
+ const translatedText = result.translatedText || result;
216
+
217
+ // Log request
218
  metrics.requests.push({
219
  timestamp: Date.now(),
220
  ip,
 
224
  inputText: text.substring(0, 100),
225
  outputText: translatedText.substring(0, 100)
226
  });
227
+
228
  res.json({ translatedText, source, target, duration });
229
  } catch (error) {
230
+ logger.error('Translation error', error);
231
  res.status(500).json({ error: 'Erreur de traduction' });
232
  }
233
  });
234
 
235
+ function translateText(text, source, target) {
236
+ return translate(text, { from: source, to: target }).then(res => res.text);
237
+ }
238
+
239
+ function findAvailableWorker() {
240
+ for (const [id, worker] of metrics.workers) {
241
+ if (worker.jobsActive < worker.maxJobs && worker.ws.readyState === 1) {
242
+ return worker;
243
+ }
244
+ }
245
+ return null;
246
+ }
247
+
248
+ function dispatchToWorker(worker, job) {
249
+ return new Promise((resolve, reject) => {
250
+ const jobId = uuidv4();
251
+ pendingJobs.set(jobId, { resolve, reject, timeout: setTimeout(() => {
252
+ pendingJobs.delete(jobId);
253
+ reject(new Error('Worker timeout'));
254
+ }, 30000) });
255
+
256
+ worker.ws.send(JSON.stringify({
257
+ type: 'translate',
258
+ jobId,
259
+ ...job
260
+ }));
261
+ worker.jobsActive++;
262
+ });
263
+ }
264
+
265
+ // Admin dashboard routes
266
+ app.get('/admin', adminAuth, (req, res) => {
267
  res.sendFile(path.join(__dirname, '../../public/admin.html'));
268
  });
269
 
 
279
  lastHeartbeat: w.lastHeartbeat
280
  });
281
  }
282
+
 
 
 
 
283
  res.json({
284
  workers,
285
  cpuHistory: metrics.cpuHistory,
286
  ramHistory: metrics.ramHistory,
287
+ requests: metrics.requests.slice(-200), // Last 200 requests
288
+ requestRate: calculateRequestRate()
 
 
 
289
  });
290
  });
291
 
292
  app.get('/api/admin/logs', adminAuth, (req, res) => {
293
+ // Return recent logs from metrics.requests
294
  const logs = metrics.requests.slice(-100).map(r => ({
295
  timestamp: r.timestamp,
296
  ip: r.ip,
 
305
 
306
  function calculateRequestRate() {
307
  const now = Date.now();
308
+ const windowMs = 120000; // 2 minutes
309
+ const recent = metrics.requests.filter(r => r.timestamp > now - windowMs);
310
+ return recent.length;
311
  }
312
 
313
+ // Main page
314
+ app.get('/', (req, res) => {
315
+ res.sendFile(path.join(__dirname, '../../public/index.html'));
316
+ });
 
317
 
318
+ // Wiki page
319
+ app.get('/wiki', (req, res) => {
320
+ res.sendFile(path.join(__dirname, '../../public/wiki.html'));
321
+ });
 
 
 
 
 
 
 
322
 
323
+ // Start server
 
 
324
  const PORT = process.env.PORT || settings.port;
325
  server.listen(PORT, '0.0.0.0', () => {
326
+ logger.info(`Central server running on port ${PORT}`);
327
+ logger.info(`Admin panel: http://localhost:${PORT}/admin`);
328
+ logger.info(`API: http://localhost:${PORT}/translate`);
329
  });