mirror of
https://github.com/jambonz/sbc-inbound.git
synced 2026-07-24 04:41:53 +00:00
dynamically recognize feature servers through their options pings
This commit is contained in:
@@ -5,8 +5,6 @@ assert.ok(process.env.JAMBONES_MYSQL_HOST &&
|
||||
process.env.JAMBONES_MYSQL_DATABASE, 'missing JAMBONES_MYSQL_XXX env vars');
|
||||
assert.ok(process.env.DRACHTIO_PORT || process.env.DRACHTIO_HOST, 'missing DRACHTIO_PORT env var');
|
||||
assert.ok(process.env.DRACHTIO_SECRET, 'missing DRACHTIO_SECRET env var');
|
||||
assert.ok(process.env.JAMBONES_RTPENGINES, 'missing JAMBONES_RTPENGINES env var');
|
||||
assert.ok(process.env.JAMBONES_FEATURE_SERVERS, 'missing JAMBONES_FEATURE_SERVERS env var');
|
||||
|
||||
const Srf = require('drachtio-srf');
|
||||
const srf = new Srf();
|
||||
@@ -25,20 +23,6 @@ const {
|
||||
connectionLimit: process.env.JAMBONES_MYSQL_CONNECTION_LIMIT || 10
|
||||
}, logger);
|
||||
|
||||
// parse rtpengines
|
||||
srf.locals.rtpEngines = process.env.JAMBONES_RTPENGINES
|
||||
.split(',')
|
||||
.map((hp) => {
|
||||
const arr = /^(.*):(.*)$/.exec(hp.trim());
|
||||
if (arr) return {host: arr[1], port: parseInt(arr[2])};
|
||||
});
|
||||
assert.ok(srf.locals.rtpEngines.length > 0, 'JAMBONES_RTPENGINES must be an array host:port addresses');
|
||||
|
||||
// parse application servers
|
||||
srf.locals.featureServers = process.env.JAMBONES_FEATURE_SERVERS
|
||||
.split(',')
|
||||
.map((hp) => hp.trim());
|
||||
|
||||
srf.locals.dbHelpers = {
|
||||
lookupAuthHook,
|
||||
lookupSipGatewayBySignalingAddress
|
||||
|
||||
+11
-6
@@ -12,10 +12,11 @@ class CallSession extends Emitter {
|
||||
this.res = res;
|
||||
this.srf = req.srf;
|
||||
this.logger = logger.child({callId: req.get('Call-ID')});
|
||||
|
||||
this.getFeatureServer = require('./fs-tracking')(this.srf, this.logger);
|
||||
}
|
||||
|
||||
async connect() {
|
||||
debug(`getRTPENGINE ${typeof getRtpEngine}`);
|
||||
const engine = getRtpEngine(this.logger);
|
||||
if (!engine) {
|
||||
this.logger.info('No available rtpengines, rejecting call!');
|
||||
@@ -27,27 +28,31 @@ class CallSession extends Emitter {
|
||||
this.answer = answer;
|
||||
this.del = del;
|
||||
|
||||
const featureServer = this.getFeatureServer();
|
||||
if (!featureServer) {
|
||||
this.logger.info('No available feature servers, rejecting call!');
|
||||
return this.res.send(480);
|
||||
}
|
||||
debug(`using feature server ${featureServer}`);
|
||||
|
||||
this.rtpEngineOpts = makeRtpEngineOpts(this.req, isWSS(this.req), false);
|
||||
this.rtpEngineResource = {destroy: this.del.bind(null, this.rtpEngineOpts.common)};
|
||||
const obj = parseUri(this.req.uri);
|
||||
const appServer = getAppserver(this.srf);
|
||||
let proxy, host, uri;
|
||||
|
||||
// replace host part of uri if its an ipv4 address, leave it otherwise
|
||||
if (/\d{1-3}\.\d{1-3}\.\d{1-3}\.\d{1-3}/.test(obj.host)) {
|
||||
host = obj.host;
|
||||
proxy = appServer;
|
||||
proxy = featureServer;
|
||||
}
|
||||
else {
|
||||
host = appServer;
|
||||
host = featureServer;
|
||||
}
|
||||
if (obj.user) uri = `${obj.scheme}:${obj.user}@${host}`;
|
||||
else uri = `${obj.scheme}:${host}`;
|
||||
debug(`uri will be: ${uri}, proxy ${proxy}`);
|
||||
|
||||
try {
|
||||
// rtpengine 'offer'
|
||||
debug('sending offer command to rtpengine');
|
||||
const response = await this.offer(this.rtpEngineOpts.offer);
|
||||
debug(`response from rtpengine to offer ${JSON.stringify(response)}`);
|
||||
if ('ok' !== response.result) {
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
const contacts = new Map();
|
||||
const noopLogger = {info: () => {}, error: () => {}};
|
||||
const debug = require('debug')('jambonz:sbc-inbound');
|
||||
|
||||
module.exports = (srf, logger) => {
|
||||
logger = logger || noopLogger;
|
||||
let dynamic = true;
|
||||
let idx = 0;
|
||||
|
||||
srf.options((req, res) => {
|
||||
res.send(200);
|
||||
if (req.has('X-FS-Status')) {
|
||||
const uri = `sip:${req.source_address}:${req.source_address}`;
|
||||
const status = req.get('X-FS-Status');
|
||||
const calls = req.has('X-FS-Calls') ? parseInt(req.get('X-FS-Calls')) : 0;
|
||||
if (status === 'open') {
|
||||
if (!contacts) {
|
||||
logger.info(`adding feature server at ${uri}`);
|
||||
//stats.gauge('sbc.featureservers', contacts.size + 1);
|
||||
}
|
||||
logger.debug(`Feature server at ${uri} has ${calls} calls`);
|
||||
contacts.set(uri, {pingTime: new Date(), calls: calls});
|
||||
}
|
||||
else {
|
||||
if (contacts.includes(uri)) {
|
||||
logger.info(`removing feature server at ${uri}`);
|
||||
contacts.delete(uri);
|
||||
//stats.gauge('sbc.featureservers', contacts.size + 1);
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
if (process.env.JAMBONES_FEATURE_SERVERS) {
|
||||
dynamic = false;
|
||||
process.env.JAMBONES_FEATURE_SERVERS
|
||||
.split(',')
|
||||
.map((hp) => hp.trim())
|
||||
.forEach((uri) => contacts.set(uri, {active: 0, calls: 0}));
|
||||
debug(`using static list of feature servers: ${[ ...contacts]}`);
|
||||
}
|
||||
|
||||
if (dynamic) {
|
||||
const CHECK_INTERVAL = 35;
|
||||
setInterval(() => {
|
||||
const dead = [];
|
||||
const deadline = (new Date.now()).getTime() - 90000;
|
||||
for (const obj of contacts) {
|
||||
if (obj[1].pingTime.getTime() < deadline) dead.push(obj[0]);
|
||||
}
|
||||
dead.forEach((uri) => {
|
||||
logger.info(`removing feature server at ${uri} due to lack of OPTIONS ping`);
|
||||
contacts.delete(uri);
|
||||
});
|
||||
|
||||
const keys = [ ...contacts.keys() ];
|
||||
logger.debug({keys}, `there are ${keys.length} feature servers online`);
|
||||
//stats.gauge('sbc.featureservers', contacts.size);
|
||||
}, CHECK_INTERVAL * 1000);
|
||||
}
|
||||
|
||||
return () => {
|
||||
let selectedUri;
|
||||
|
||||
if (dynamic) {
|
||||
const featureServers = [ ...contacts].map((o) => Object.assign({}, {uri: o[0]}, o[1]));
|
||||
debug({featureServers}, 'selecting feature servers with least calls');
|
||||
const fs = featureServers.sort((a, b) => (a.calls - b.calls)).shift();
|
||||
if (!fs) logger.info('No available feature servers!');
|
||||
else selectedUri = fs.uri;
|
||||
}
|
||||
else {
|
||||
const featureServers = [ ...contacts].map((o) => o[0]);
|
||||
debug({featureServers}, 'selecting feature servers from static list');
|
||||
selectedUri = featureServers[idx++ % featureServers.length];
|
||||
logger.debug(`selected ${selectedUri}`);
|
||||
}
|
||||
return selectedUri;
|
||||
};
|
||||
};
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "sbc-inbound",
|
||||
"version": "0.2.1",
|
||||
"version": "0.3.0",
|
||||
"main": "app.js",
|
||||
"engines": {
|
||||
"node": ">= 10.16.0"
|
||||
|
||||
Reference in New Issue
Block a user