revamp to mysql for gateway configuration and major code refactor

This commit is contained in:
Dave Horton committed 2019-12-13 10:20:27 -05:00
1 parent c22b6f9c08
commit 5a529b2bb9
23 files changed
+1152 -187

No files matched your search

+131
View File
@@ -0,0 +1,131 @@
const Emitter = require('events');
const config = require('config');
const Client = require('rtpengine-client').Client ;
const rtpengine = new Client();
const offer = rtpengine.offer.bind(rtpengine, config.get('rtpengine'));
const answer = rtpengine.answer.bind(rtpengine, config.get('rtpengine'));
const del = rtpengine.delete.bind(rtpengine, config.get('rtpengine'));
const {getAppserver, isWSS, makeRtpEngineOpts} = require('./utils');
const {forwardInDialogRequests} = require('drachtio-fn-b2b-sugar');
const {parseUri, SipError} = require('drachtio-srf');
const debug = require('debug')('jambonz:sbc-inbound');
class CallSession extends Emitter {
constructor(logger, req, res) {
super();
this.req = req;
this.res = res;
this.srf = req.srf;
this.logger = logger.child({callId: req.get('Call-ID')});
}
async connect() {
this.rtpEngineOpts = makeRtpEngineOpts(this.req, isWSS(this.req), false);
this.rtpEngineResource = {destroy: del.bind(rtpengine, this.rtpEngineOpts.common)};
const obj = parseUri(this.req.uri);
const appServer = getAppserver();
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;
}
else {
host = appServer;
}
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 offer(this.rtpEngineOpts.offer);
debug(`response from rtpengine to offer ${JSON.stringify(response)}`);
if ('ok' !== response.result) {
this.logger.error(`rtpengine offer failed with ${JSON.stringify(response)}`);
throw new Error('rtpengine failed: answer');
}
// now send the INVITE in towards the feature servers
const headers = {'X-Forwarded-For': this.req.source_address};
if (this.req.locals.carrier) Object.assign(headers, {'X-Originating-Carrier': this.req.locals.carrier});
debug(`sending INVITE to ${proxy} with ${uri}`);
const {uas, uac} = await this.srf.createB2BUA(this.req, this.res, uri, {
proxy,
headers,
proxyRequestHeaders: ['all'],
proxyResponseHeaders: ['all'],
localSdpB: response.sdp,
localSdpA: async(sdp, res) => {
const opts = Object.assign({sdp, 'to-tag': res.getParsedHeader('To').params.tag},
this.rtpEngineOpts.answer);
const response = await answer(opts);
if ('ok' !== response.result) {
this.logger.error(`rtpengine answer failed with ${JSON.stringify(response)}`);
throw new Error('rtpengine failed: answer');
}
return response.sdp;
}
});
// successfully connected
this.logger.info('call connected');
debug('call connected');
this.emit('connected');
this._setHandlers({uas, uac});
return;
} catch (err) {
if (err instanceof SipError) {
this.logger.info(`call failed with ${err.status}`);
this.emit('failed');
this.rtpEngineResource.destroy();
}
}
}
_setHandlers({uas, uac}) {
this.uas = uas;
this.uac = uac;
[uas, uac].forEach((dlg) => {
//hangup
dlg.on('destroy', () => {
this.logger.info('call ended');
this.rtpEngineResource.destroy();
});
//re-invite
dlg.on('modify', this._onReinvite.bind(this, dlg));
});
// default forwarding of other request types
forwardInDialogRequests(uas);
}
async _onReinvite(dlg, req, res) {
try {
let response = await offer(Object.assign({sdp: req.body}, this.rtpEngineOpts.offer));
if ('ok' !== response.result) {
res.send(488);
throw new Error(`_onReinvite: rtpengine failed: offer: ${JSON.stringify(response)}`);
}
const sdp = await dlg.other.modify(response.sdp);
const opts = Object.assign({sdp, 'to-tag': res.getParsedHeader('To').params.tag},
this.rtpEngineOpts.answer);
response = await answer(opts);
if ('ok' !== response.result) {
res.send(488);
throw new Error(`_onReinvite: rtpengine failed: ${JSON.stringify(response)}`);
}
res.send(200, {body: response.sdp});
} catch (err) {
this.logger.error(err, 'Error handling reinvite');
}
}
}
module.exports = CallSession;
-62
View File
@@ -1,62 +0,0 @@
const config = require('config');
const Client = require('rtpengine-client').Client ;
const rtpengine = new Client();
const offer = rtpengine.offer.bind(rtpengine, config.get('rtpengine'));
const answer = rtpengine.answer.bind(rtpengine, config.get('rtpengine'));
const del = rtpengine.delete.bind(rtpengine, config.get('rtpengine'));
const {getAppserver, isWSS, makeRtpEngineOpts} = require('./utils');
module.exports = handler;
function handler({log}) {
return async(req, res) => {
const logger = log.child({callId: req.get('Call-ID')});
const srf = req.srf;
const rtpEngineOpts = makeRtpEngineOpts(req, isWSS(req), false);
const rtpEngineResource = {destroy: del.bind(rtpengine, rtpEngineOpts.common)};
const uri = getAppserver();
logger.info(`received inbound INVITE from ${req.protocol}/${req.source_address}:${req.source_port}`);
try {
const response = await offer(rtpEngineOpts.offer);
if ('ok' !== response.result) {
res.send(480);
throw new Error(`failed allocating rtpengine endpoint: ${JSON.stringify(response)}`);
}
const {uas, uac} = await srf.createB2BUA(req, res, uri, {
headers: {
'X-Forwarded-For': req.source_address,
'X-Forwarded-Proto': req.getParsedHeader('Via')[0].protocol.toLowerCase(),
'X-Forwarded-Carrier': req.carrier_name
},
proxyRequestHeaders: ['User-Agent', 'Subject'],
localSdpB: response.sdp,
localSdpA: (sdp, res) => {
const opts = Object.assign({sdp, 'to-tag': res.getParsedHeader('To').params.tag},
rtpEngineOpts.answer);
return answer(opts)
.then((response) => {
if ('ok' !== response.result) throw new Error('error allocating rtpengine');
return response.sdp;
});
}
});
logger.info('call connected');
setHandlers(logger, uas, uac, rtpEngineResource);
} catch (err) {
logger.error(err, 'Error connecting call');
rtpEngineResource.destroy();
}
};
}
function setHandlers(logger, uas, uac, rtpEngineResource) {
[uas, uac].forEach((dlg) => dlg.on('destroy', () => {
logger.info('call ended');
dlg.other.destroy();
rtpEngineResource.destroy();
}));
//TODO: handle re-INVITEs, REFER, INFO
}
+25 -9
View File
@@ -1,12 +1,28 @@
const {fromInboundTrunk} = require('./utils');
const config = require('config');
const authenticator = require('drachtio-http-authenticator')(config.get('authCallback'));
const debug = require('debug')('jambonz:sbc-inbound');
function auth(req, res, next) {
if (fromInboundTrunk(req)) {
return next();
module.exports = function(srf, logger) {
const {lookupSipGatewayBySignalingAddress, lookupAuthHook} = srf.locals.dbHelpers;
const authenticator = require('drachtio-http-authenticator')(lookupAuthHook, logger);
async function challengeDeviceCalls(req, res, next) {
req.locals = req.locals || {};
try {
const gateway = await lookupSipGatewayBySignalingAddress(req.source_address, req.source_port);
if (!gateway) {
req.locals.originator = 'device';
return authenticator(req, res, next);
}
debug(`challengeDeviceCalls: call came from gateway: ${JSON.stringify(gateway)}`);
req.locals.originator = 'trunk';
req.locals.carrier = gateway.name;
next();
} catch (err) {
logger.error(err, `${req.get('Call-ID')} Error looking up related info for inbound call`);
res.send(500);
}
}
authenticator(req, res, next);
}
module.exports = { auth };
return {
challengeDeviceCalls
};
};
-11
View File
@@ -1,16 +1,6 @@
const config = require('config');
let idx = 0;
function fromInboundTrunk(req) {
const trunks = config.has('trunks.inbound') ?
config.get('trunks.inbound') : [];
if (isWSS(req)) return false;
const trunk = trunks.find((t) => t.host.includes(req.source_address));
if (!trunk) return false;
req.carrier_name = trunk.name;
return true;
}
function isWSS(req) {
return req.getParsedHeader('Via')[0].protocol.toLowerCase().startsWith('ws');
}
@@ -34,7 +24,6 @@ function makeRtpEngineOpts(req, srcIsUsingSrtp, dstIsUsingSrtp) {
}
module.exports = {
fromInboundTrunk,
isWSS,
getAppserver,
makeRtpEngineOpts