diff --git a/app.js b/app.js index 149f6ce..f9c2adf 100644 --- a/app.js +++ b/app.js @@ -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 diff --git a/lib/call-session.js b/lib/call-session.js index 35d2a78..0201b04 100644 --- a/lib/call-session.js +++ b/lib/call-session.js @@ -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) { diff --git a/lib/fs-tracking.js b/lib/fs-tracking.js new file mode 100644 index 0000000..4cf74d5 --- /dev/null +++ b/lib/fs-tracking.js @@ -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; + }; +}; diff --git a/package.json b/package.json index b692424..962893e 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "sbc-inbound", - "version": "0.2.1", + "version": "0.3.0", "main": "app.js", "engines": { "node": ">= 10.16.0"