obcon-scada / app /include /cluster.js
chanmin0723's picture
Initial obcon SCADA deploy
e4bf523
Raw
History Blame Contribute Delete
7.07 kB
'use strict'
/**
* Copyright (c) 2017~2019, OBCon Inc.
* All rights reserved.
*/
/**
* @file
* @copyright 2017~2019, OBCon Inc.
* @author gye hyun james kim [pnuskgh@gmail.com]
*/
let os = require('os');
let cluster = require('cluster');
let moment = require('moment');
let { v1: uuidv1 } = require('uuid');
// ✅ 누락돼서 터졌던 caches 정의
const caches = require('./cache');
/**
* 안전하게 worker.send 하기 (죽은 worker/닫힌 IPC 채널 방지)
*/
function safeSend(worker, msg) {
try {
if (!worker) return false;
// node 버전에 따라 worker.isConnected()가 안전
if (typeof worker.isConnected === 'function' && !worker.isConnected()) return false;
if (worker.process && worker.process.connected === false) return false;
worker.send(msg);
return true;
} catch (e) {
// ERR_IPC_CHANNEL_CLOSED / EPIPE는 “죽은 채널에 보냈다”라서 무시해도 됨
return false;
}
}
class Cluster {
constructor() {
// cluster.schedulingPolicy = cluster.SCHED_NONE; //--- 워커 스케쥴을 OS에 맡긴다.
cluster.schedulingPolicy = cluster.SCHED_RR; //---워커 스케쥴을 Round Robin 방식으로 한다.
this.init_event();
}
get cpus() {
return os.cpus().length;
}
init_event() {
if (cluster.isMaster) {
this.init_event_master();
} else {
this.init_event_worker();
}
}
init_event_master() {
//--- master:setup > master:fork > master:online
//--- worker:disconnect > master:disconnect > master:exit
cluster.on("setup", function(mySettings){
// console.log("Master :: Master %s is setup", JSON.stringify(mySettings));
});
cluster.on("fork", function(worker) {
// console.log("Master :: Worker %s is fork", worker.id);
});
cluster.on("online", function(worker) {
console.log("Master :: Worker %s is online", worker.id);
});
cluster.on("disconnect", function(worker) {
console.log("Master :: Worker %s is disconnect", worker.id);
});
cluster.on("exit", function(worker, code, signal) {
console.log("Master :: Worker %s is exit", worker.id);
console.log("Master :: Worker %s is exit code : %s", worker.id, code);
console.log("Master :: Worker %s is exit signal : %s", worker.id, signal);
});
cluster.on("listening", function(worker, address) {
console.log("Master :: Worker %s is listening", worker.id);
});
//--- Master에게 온 메시지를 모든 worker에게 전달한다.
// ✅ 기존 코드: 죽은 worker에게도 send → EPIPE/IPC_CLOSED 폭탄
cluster.on('message', function(_workerSender, req) {
for (const idx in cluster.workers) {
safeSend(cluster.workers[idx], req);
}
});
}
init_event_worker() {
//--- worker:disconnect > master:disconnect > master:exit
cluster.worker.on("disconnect", function() {
console.log("Worker :: Worker %s disconnect", cluster.worker.id);
});
cluster.worker.on("error", function(error) {
console.log("Worker :: Worker %s error", cluster.worker.id);
});
cluster.worker.on("exit", function(code, signal) {
console.log("Worker :: Worker %s exit", cluster.worker.id);
});
cluster.worker.on("listening", function(address) {
console.log("Worker :: Worker %s listening", cluster.worker.id);
});
//--- Cache에 저장된 handler에게 메시지 처리를 위임 한다.
cluster.worker.on("message", function(req) {
if (typeof(req && req.headers) !== 'undefined') {
if ((req.headers.module == 'cluster') && (req.headers.action == 'workerInit')) {
global.appl.worker = req.body;
}
} else {
const handler = caches.getCache(req && req.type, req && req.path);
if (handler == null) {
console.log("%s - %s handler가 없습니다.", req && req.type, req && req.path);
} else {
console.log("%s - %s 요청 처리", req.type, req.path);
handler(req);
}
}
});
cluster.worker.on("online", function() {
console.log("Worker :: Worker %s online", cluster.worker.id);
});
}
/**
* config.cluster 설정 기반으로 worker fork
* (PM2가 아닌 직접 실행 모드에서 사용)
*/
forks() {
if (!cluster.isMaster) return;
const countHttp = Number((global.config && global.config.cluster && global.config.cluster.worker_count_http) || 0);
const countTcp = Number((global.config && global.config.cluster && global.config.cluster.worker_count_tcp) || 0);
const countUdp = Number((global.config && global.config.cluster && global.config.cluster.worker_count_udp) || 0);
const countModbus = Number((global.config && global.config.cluster && global.config.cluster.worker_count_modbus) || 0);
const countProxy = Number((global.config && global.config.cluster && global.config.cluster.worker_count_proxy) || 0);
let total = 0;
if (global.appl && global.appl.type === 'obcon_proxy') {
total = Math.max(1, countProxy || 1);
} else {
total = countHttp + countTcp + countUdp + countModbus;
if (total <= 0) total = 1;
}
for (let seq = 1; seq <= total; seq++) {
const worker = cluster.fork();
const workerType = this._getWorkerType(seq, {
countHttp,
countTcp,
countUdp,
countModbus,
countProxy
});
const workerInfo = {
id: worker.id,
type: workerType,
seq
};
worker.on('online', () => {
safeSend(worker, {
headers: { module: 'cluster', action: 'workerInit' },
body: workerInfo,
ts: moment().format('YYYY-MM-DD HH:mm:ss.SSS'),
uuid: uuidv1()
});
});
}
}
_getWorkerType(seq, counts) {
if (global.appl && global.appl.type === 'obcon_proxy') return 'proxy';
if (seq <= counts.countHttp) return 'http';
if (seq <= counts.countHttp + counts.countTcp) return 'tcp';
if (seq <= counts.countHttp + counts.countTcp + counts.countUdp) return 'udp';
if (seq <= counts.countHttp + counts.countTcp + counts.countUdp + counts.countModbus) return 'modbus';
return 'http';
}
/**
* Worker에 들어온 요청(type/path)을 처리할 handler 등록
*/
setHandler(type, path, handler) {
caches.setCache(type, path, handler);
}
/**
* workerInit 브로드캐스트
*/
sendWorkerInit(workerInfo) {
if (cluster.isMaster) {
const msg = {
headers: { module: 'cluster', action: 'workerInit' },
body: workerInfo,
ts: moment().format('YYYY-MM-DD HH:mm:ss.SSS'),
uuid: uuidv1()
};
for (const idx in cluster.workers) {
safeSend(cluster.workers[idx], msg);
}
}
}
/**
* 임의 메시지 브로드캐스트 (마스터만)
*/
broadcast(req) {
if (!cluster.isMaster) return;
for (const idx in cluster.workers) {
safeSend(cluster.workers[idx], req);
}
}
}
module.exports = new Cluster();