Compare commits

..
Author SHA1 Message Date
Dave Horton 24c2225758 call rtpengine offer using callback 2022-09-27 16:03:46 +01:00
9 changed files with 484 additions and 1445 deletions
+1 -1
View File
@@ -1,4 +1,4 @@
FROM --platform=linux/amd64 node:18.9.0-alpine3.16 as base
FROM --platform=linux/amd64 node:18.8.0-alpine as base
RUN apk --update --no-cache add --virtual .builds-deps build-base python3
+1 -1
View File
@@ -69,7 +69,7 @@ const {
const ngProtocol = process.env.JAMBONES_NG_PROTOCOL || 'udp';
const ngPort = process.env.RTPENGINE_PORT || ('udp' === ngProtocol ? 22222 : 8080);
const {getRtpEngine, setRtpEngines} = require('@jambonz/rtpengine-utils')([], logger, {
//emitter: stats,
emitter: stats,
dtmfListenPort: process.env.DTMF_LISTEN_PORT || 22224,
protocol: ngProtocol
});
+6 -28
View File
@@ -1,7 +1,7 @@
const Emitter = require('events');
const bent = require('bent');
const assert = require('assert');
const PORT = process.env.AWS_SNS_PORT || 3010;
const PORT = process.env.AWS_SNS_PORT || 3001;
const {LifeCycleEvents} = require('./constants');
const express = require('express');
const app = express();
@@ -21,26 +21,6 @@ class SnsNotifier extends Emitter {
this.logger = logger;
}
_doListen(logger, app, port, resolve) {
return app.listen(port, () => {
this.snsEndpoint = `http://${this.publicIp}:${port}`;
logger.info(`SNS lifecycle server listening on http://localhost:${port}`);
resolve(app);
});
}
_handleErrors(logger, app, resolve, reject, e) {
if (e.code === 'EADDRINUSE' &&
process.env.AWS_SNS_PORT_MAX &&
e.port < process.env.AWS_SNS_PORT_MAX) {
logger.info(`SNS lifecycle server failed to bind port on ${e.port}, will try next port`);
const server = this._doListen(logger, app, ++e.port, resolve);
server.on('error', this._handleErrors.bind(this, logger, app, resolve, reject));
return;
}
reject(e);
}
async _handlePost(req, res) {
try {
@@ -65,7 +45,6 @@ class SnsNotifier extends Emitter {
}, 'response from SNS SubscribeURL');
const data = await this.describeInstance();
this.lifecycleState = data.AutoScalingInstances[0].LifecycleState;
this.emit('SubscriptionConfirmation', {publicIp: this.publicIp});
break;
case 'Notification':
@@ -101,12 +80,14 @@ class SnsNotifier extends Emitter {
async init() {
try {
this.logger.debug('SnsNotifier: retrieving instance data');
this.logger.info('SnsNotifier: retrieving instance data');
this.instanceId = await getString('http://169.254.169.254/latest/meta-data/instance-id');
this.publicIp = await getString('http://169.254.169.254/latest/meta-data/public-ipv4');
this.snsEndpoint = `http://${this.publicIp}:${PORT}`;
this.logger.info({
instanceId: this.instanceId,
publicIp: this.publicIp
publicIp: this.publicIp,
snsEndpoint: this.snsEndpoint
}, 'retrieved AWS instance data');
// start listening
@@ -118,10 +99,7 @@ class SnsNotifier extends Emitter {
this.logger.error(err, 'burped error');
res.status(err.status || 500).json({msg: err.message});
});
return new Promise((resolve, reject) => {
const server = this._doListen(this.logger, app, PORT, resolve);
server.on('error', this._handleErrors.bind(this, this.logger, app, resolve, reject));
});
app.listen(PORT);
} catch (err) {
this.logger.error({err}, 'Error retrieving AWS instance metadata');
+136 -165
View File
@@ -80,28 +80,6 @@ class CallSession extends Emitter {
return tp && -1 !== tp.indexOf('SAVP');
}
get isFive9VoiceStream() {
return this.req.has('X-Five9-StreamingPairId');
}
get isPossibleWebRtcClient() {
return this.req.locals.isPossibleWebRtcClient;
}
subscribeForDTMF(dlg) {
if (!this._subscribedForDTMF) {
this._subscribedForDTMF = true;
this.subscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag,
this._onDTMF.bind(this, dlg));
}
}
unsubscribeForDTMF() {
if (this._subscribedForDTMF) {
this._subscribedForDTMF = false;
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
}
}
async connect() {
const {sdp} = this.req.locals;
this.logger.info('inbound call accepted for routing');
@@ -151,8 +129,7 @@ class CallSession extends Emitter {
}
this.logger.debug(`using feature server ${featureServer}`);
const wantsSrtp = this.req.locals.possibleWebRtcClient = SdpWantsSrtp(sdp);
this.rtpEngineOpts = makeRtpEngineOpts(this.req, wantsSrtp, false, this.isFromMSTeams);
this.rtpEngineOpts = makeRtpEngineOpts(this.req, SdpWantsSrtp(sdp), false, this.isFromMSTeams);
this.rtpEngineResource = {destroy: this.del.bind(null, this.rtpEngineOpts.common)};
const obj = parseUri(this.req.uri);
let proxy, host, uri;
@@ -171,145 +148,144 @@ class CallSession extends Emitter {
else uri = `${obj.scheme}:${host}`;
this.logger.info(`uri will be: ${uri}, proxy ${proxy}`);
try {
const opts = {
...this.rtpEngineOpts.common,
...this.rtpEngineOpts.uac.mediaOpts,
'from-tag': this.rtpEngineOpts.uas.tag,
direction: ['public', 'private'],
sdp
};
const startAt = process.hrtime();
const response = await this.offer(opts);
this.rtpengineIp = opts.sdp ? parseConnectionIp(opts.sdp) : 'undefined';
const rtt = roundTripTime(startAt);
this.stats.histogram('app.rtpengine.response_time', rtt, [
'direction:inbound', 'command:offer', `rtpengine:${this.rtpengineIp}`]);
this.logger.debug({opts, response, rtt, rtpengine: this.rtpengineIp}, 'response from rtpengine to offer');
if ('ok' !== response.result) {
this.logger.error({}, `rtpengine offer failed with ${JSON.stringify(response)}`);
throw new Error('rtpengine failed: answer');
}
let headers = {
'From': createBLegFromHeader(this.req),
'To': this.req.get('To'),
'X-Account-Sid': this.req.locals.account_sid,
'X-CID': this.req.get('Call-ID'),
'X-Forwarded-For': `${this.req.source_address}`
};
if (this.privateSipAddress) headers = {...headers, Contact: `<sip:${this.privateSipAddress}>`};
const spdOfferB = this.siprec && this.xml ?
createSiprecBody(headers, response.sdp, this.xml.type, this.xml.content) :
response.sdp;
const responseHeaders = {};
if (this.req.locals.carrier) {
Object.assign(headers, {
'X-Originating-Carrier': this.req.locals.carrier,
'X-Voip-Carrier-Sid': this.req.locals.voip_carrier_sid
});
}
if (this.req.locals.msTeamsTenantFqdn) {
Object.assign(headers, {'X-MS-Teams-Tenant-FQDN': this.req.locals.msTeamsTenantFqdn});
// for Microsoft Teams the Contact header must include the tenant FQDN
Object.assign(responseHeaders, {
Allow: 'INVITE, ACK, OPTIONS, CANCEL, BYE, NOTIFY, UPDATE, PRACK',
Contact: `sip:${this.req.locals.msTeamsTenantFqdn}`
});
}
if (this.req.locals.application_sid) {
Object.assign(headers, {'X-Application-Sid': this.req.locals.application_sid});
}
else if (this.req.authorization) {
if (this.req.authorization.grant && this.req.authorization.grant.application_sid) {
Object.assign(headers, {'X-Application-Sid': this.req.authorization.grant.application_sid});
const opts = {
...this.rtpEngineOpts.common,
...this.rtpEngineOpts.uac.mediaOpts,
'from-tag': this.rtpEngineOpts.uas.tag,
direction: ['public', 'private'],
sdp
};
const startAt = process.hrtime();
this.offer(opts, async(err, response) => {
try {
if (err || 'ok' !== response.result) {
this.logger.error({}, `rtpengine offer failed with ${JSON.stringify(response)}`);
return this.res.send(480, {
headers: {
'X-Reason': 'RTPENGINE_OFFER_FAILED'
}
});
}
if (this.req.authorization.challengeResponse) {
const {username, realm} = this.req.authorization.challengeResponse;
Object.assign(headers, {'X-Authenticated-User': `${username}@${realm}`});
this.rtpengineIp = opts.sdp ? parseConnectionIp(opts.sdp) : 'undefined';
const rtt = roundTripTime(startAt);
this.stats.histogram('app.rtpengine.response_time', rtt, [
'direction:inbound', 'command:offer', `rtpengine:${this.rtpengineIp}`]);
this.logger.debug({opts, response, rtt, rtpengine: this.rtpengineIp}, 'response from rtpengine to offer');
if ('ok' !== response.result) {
}
}
let headers = {
'From': createBLegFromHeader(this.req),
'To': this.req.get('To'),
'X-Account-Sid': this.req.locals.account_sid,
'X-CID': this.req.get('Call-ID'),
'X-Forwarded-For': `${this.req.source_address}`
};
if (this.privateSipAddress) headers = {...headers, Contact: `<sip:${this.privateSipAddress}>`};
if (this.req.canceled) throw new Error('call canceled');
const spdOfferB = this.siprec && this.xml ?
createSiprecBody(headers, response.sdp, this.xml.type, this.xml.content) :
response.sdp;
// now send the INVITE in towards the feature servers
debug(`sending INVITE to ${proxy} with ${uri}`);
const {uas, uac} = await this.srf.createB2BUA(this.req, this.res, uri, {
proxy,
headers,
responseHeaders,
proxyRequestHeaders: [
'all',
'-Authorization',
'-Max-Forwards',
'-Record-Route',
'-Session-Expires',
'-X-Subspace-Forwarded-For'
],
proxyResponseHeaders: ['all', '-X-Trace-ID'],
localSdpB: spdOfferB,
localSdpA: async(sdp, res) => {
this.rtpEngineOpts.uac.tag = res.getParsedHeader('To').params.tag;
const opts = {
...this.rtpEngineOpts.common,
...this.rtpEngineOpts.uas.mediaOpts,
'from-tag': this.rtpEngineOpts.uas.tag,
'to-tag': this.rtpEngineOpts.uac.tag,
sdp
};
const startAt = process.hrtime();
const response = await this.answer(opts);
const rtt = roundTripTime(startAt);
this.stats.histogram('app.rtpengine.response_time', rtt, [
'direction:inbound', 'command:answer', `rtpengine:${this.rtpengineIp}`]);
if ('ok' !== response.result) {
this.logger.error(`rtpengine answer failed with ${JSON.stringify(response)}`);
throw new Error('rtpengine failed: answer');
const responseHeaders = {};
if (this.req.locals.carrier) {
Object.assign(headers, {
'X-Originating-Carrier': this.req.locals.carrier,
'X-Voip-Carrier-Sid': this.req.locals.voip_carrier_sid
});
}
if (this.req.locals.msTeamsTenantFqdn) {
Object.assign(headers, {'X-MS-Teams-Tenant-FQDN': this.req.locals.msTeamsTenantFqdn});
// for Microsoft Teams the Contact header must include the tenant FQDN
Object.assign(responseHeaders, {
Allow: 'INVITE, ACK, OPTIONS, CANCEL, BYE, NOTIFY, UPDATE, PRACK',
Contact: `sip:${this.req.locals.msTeamsTenantFqdn}`
});
}
if (this.req.locals.application_sid) {
Object.assign(headers, {'X-Application-Sid': this.req.locals.application_sid});
}
else if (this.req.authorization) {
if (this.req.authorization.grant && this.req.authorization.grant.application_sid) {
Object.assign(headers, {'X-Application-Sid': this.req.authorization.grant.application_sid});
}
/* special case: Five9 Voicestream calls do not advertise a:sendonly, though they should */
if (this.isFive9VoiceStream) {
if (this.req.authorization.challengeResponse) {
const {username, realm} = this.req.authorization.challengeResponse;
Object.assign(headers, {'X-Authenticated-User': `${username}@${realm}`});
}
}
if (this.req.canceled) throw new Error('call canceled');
// now send the INVITE in towards the feature servers
debug(`sending INVITE to ${proxy} with ${uri}`);
const {uas, uac} = await this.srf.createB2BUA(this.req, this.res, uri, {
proxy,
headers,
responseHeaders,
proxyRequestHeaders: [
'all',
'-Authorization',
'-Max-Forwards',
'-Record-Route',
'-Session-Expires',
'-X-Subspace-Forwarded-For'
],
proxyResponseHeaders: ['all', '-X-Trace-ID'],
localSdpB: spdOfferB,
localSdpA: async(sdp, res) => {
this.rtpEngineOpts.uac.tag = res.getParsedHeader('To').params.tag;
const opts = {
...this.rtpEngineOpts.common,
'from-tag':this.rtpEngineOpts.uac.tag
...this.rtpEngineOpts.uas.mediaOpts,
'from-tag': this.rtpEngineOpts.uas.tag,
'to-tag': this.rtpEngineOpts.uac.tag,
sdp
};
this.logger.info('Voicestream call from Five9, blocking audio in the reverse direction');
const response = await Promise.all([this.blockMedia(opts), this.blockDTMF(opts)]);
this.logger.debug({response}, 'response to blockMedia/blockDTMF');
const startAt = process.hrtime();
const response = await this.answer(opts);
const rtt = roundTripTime(startAt);
this.stats.histogram('app.rtpengine.response_time', rtt, [
'direction:inbound', 'command:answer', `rtpengine:${this.rtpengineIp}`]);
if ('ok' !== response.result) {
this.logger.error(`rtpengine answer failed with ${JSON.stringify(response)}`);
throw new Error('rtpengine failed: answer');
}
return response.sdp;
}
});
return response.sdp;
// successfully connected
this.logger.info('call connected successfully to feature server');
debug('call connected successfully to feature server');
this._setHandlers({uas, uac});
} catch (err) {
// eslint-disable-next-line promise/no-promise-in-callback
this.rtpEngineResource.destroy()
.catch((err) => this.logger.info({err}, 'Error destroying rtpe after failure'));
this.activeCallIds.delete(this.req.get('Call-ID'));
this.stats.gauge('sbc.sip.calls.count', this.activeCallIds.size);
if (err instanceof SipError) {
const tags = ['accepted:no', `sipStatus:${err.status}`, `originator:${this.req.locals.originator}`];
this.stats.increment('sbc.terminations', tags);
this.logger.info(`call failed to connect to feature server with ${err.status}`);
this.emit('failed');
}
});
// successfully connected
this.logger.info('call connected successfully to feature server');
debug('call connected successfully to feature server');
this._setHandlers({uas, uac});
return;
} catch (err) {
this.rtpEngineResource.destroy().catch((err) => this.logger.info({err}, 'Error destroying rtpe after failure'));
this.activeCallIds.delete(this.req.get('Call-ID'));
this.stats.gauge('sbc.sip.calls.count', this.activeCallIds.size);
if (err instanceof SipError) {
const tags = ['accepted:no', `sipStatus:${err.status}`, `originator:${this.req.locals.originator}`];
this.stats.increment('sbc.terminations', tags);
this.logger.info(`call failed to connect to feature server with ${err.status}`);
this.emit('failed');
else if (err.message !== 'call canceled') {
this.logger.error(err, 'unexpected error routing inbound call');
}
this.srf.endSession(this.req);
}
else if (err.message !== 'call canceled') {
this.logger.error(err, 'unexpected error routing inbound call');
}
this.srf.endSession(this.req);
}
});
}
_setDlgHandlers(dlg) {
const {callId} = dlg.sip;
this.activeCallIds.set(callId, this);
if (this.isPossibleWebRtcClient) this.subscribeForDTMF(dlg);
this.subscribeDTMF(this.logger, callId, this.rtpEngineOpts.uas.tag,
this._onDTMF.bind(this));
dlg.on('destroy', () => {
debug('call ended with normal termination');
this.logger.info('call ended with normal termination');
@@ -357,8 +333,7 @@ class CallSession extends Emitter {
try {
await other.destroy();
} catch (err) {}
this.unsubscribeForDTMF();
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
const trackingOn = process.env.JAMBONES_TRACK_ACCOUNT_CALLS ||
process.env.JAMBONES_TRACK_SP_CALLS ||
@@ -414,7 +389,8 @@ class CallSession extends Emitter {
});
});
if (this.isPossibleWebRtcClient) this.subscribeForDTMF(uac);
this.subscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag,
this._onDTMF.bind(this, uac));
uas.on('modify', this._onReinvite.bind(this, uas));
uac.on('modify', this._onReinvite.bind(this, uac));
@@ -508,7 +484,7 @@ Duration=${payload.duration} `
});
}
this.unsubscribeForDTMF();
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
const uas = await this.srf.createUAS(req, res, {
localSdp: response.sdp,
@@ -541,16 +517,11 @@ Duration=${payload.duration} `
req.body;
const reason = req.get('X-Reason');
const isReleasingMedia = reason && dlg.type === 'uas' && ['release-media', 'anchor-media'].includes(reason);
const fromTag = dlg.type === 'uas' ? this.rtpEngineOpts.uas.tag : this.rtpEngineOpts.uac.tag;
const toTag = dlg.type === 'uas' ? this.rtpEngineOpts.uac.tag : this.rtpEngineOpts.uas.tag;
const offerMedia = dlg.type === 'uas' ? this.rtpEngineOpts.uac.mediaOpts : this.rtpEngineOpts.uas.mediaOpts;
const answerMedia = dlg.type === 'uas' ? this.rtpEngineOpts.uas.mediaOpts : this.rtpEngineOpts.uac.mediaOpts;
const direction = dlg.type === 'uas' ? ['public', 'private'] : ['private', 'public'];
if (isReleasingMedia) {
if (!offerMedia.flags.includes('asymmetric')) offerMedia.flags.push('asymmetric');
offerMedia.flags = offerMedia.flags.filter((f) => f !== 'media handover');
}
let opts = {
...this.rtpEngineOpts.common,
...offerMedia,
@@ -569,11 +540,11 @@ Duration=${payload.duration} `
/* if this is a re-invite from the FS to change media anchoring, avoid sending the reinvite out */
let sdp;
if (isReleasingMedia && !this.callerIsUsingSrtp) {
if (reason && dlg.type === 'uac' && ['release-media', 'anchor-media'].includes(reason) &&
!this.callerIsUsingSrtp) {
this.logger.info({response}, `got a reinvite from FS to ${reason}`);
sdp = dlg.other.remote.sdp;
if (!answerMedia.flags.includes('asymmetric')) answerMedia.flags.push('asymmetric');
answerMedia.flags = answerMedia.flags.filter((f) => f !== 'media handover');
answerMedia.flags = ['asymmetric', 'port latching'];
this._mediaReleased = 'release-media' === reason;
}
else {
@@ -814,15 +785,15 @@ Duration=${payload.duration} `
res.send(202);
// invite to new fs
const headers = {
...(req.has('X-Retain-Call-Sid') && {'X-Retain-Call-Sid': req.get('X-Retain-Call-Sid')}),
...(req.has('X-Account-Sid') && {'X-Account-Sid': req.get('X-Account-Sid')})
};
const headers = {};
if (req.has('X-Retain-Call-Sid')) {
Object.assign(headers, {'X-Retain-Call-Sid': req.get('X-Retain-Call-Sid')});
}
const uac = await this.srf.createUAC(referTo.uri, {localSdp: dlg.local.sdp, headers});
this.uac = uac;
uac.other = this.uas;
this.uas.other = uac;
uac.on('modify', this._onReinvite.bind(this, uac));
uac.on('modify', this._onFeatureServerReinvite.bind(this, uac));
uac.on('refer', this._onFeatureServerTransfer.bind(this, uac));
uac.on('destroy', () => {
this.logger.info('call ended with normal termination');
@@ -840,7 +811,7 @@ Duration=${payload.duration} `
const response = await this.answer(opts);
if ('ok' !== response.result) {
res.send(488);
throw new Error(`_onFeatureServerTransfer: rtpengine failed: ${JSON.stringify(response)}`);
throw new Error(`_onFeatureServerReinvite: rtpengine failed: ${JSON.stringify(response)}`);
}
this.logger.info('successfully moved call to new feature server');
} catch (err) {
@@ -901,7 +872,7 @@ Duration=${payload.duration} `
// successfully connected
this.logger.info('successfully connected new call leg for REFER');
this.unsubscribeForDTMF();
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
this.referInvite = null;
sendNotify(this.uas, '200 OK');
this.uas.destroy();
+1 -55
View File
@@ -57,16 +57,6 @@ WHERE sg.voip_carrier_sid = ?
AND sg.voip_carrier_sid = vc.voip_carrier_sid
AND outbound = 1`;
const sqlSelectCarrierRequiringRegistration = `
SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.service_provider_sid, vc.account_sid,
vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask
FROM sip_gateways sg, voip_carriers vc
WHERE sg.voip_carrier_sid = vc.voip_carrier_sid
AND vc.requires_register = 1
AND vc.is_active = 1
AND vc.register_sip_realm = ?
AND vc.register_username = ?`;
const gatewayMatchesSourceAddress = (source_address, gw) => {
if (32 === gw.netmask && gw.ipv4 === source_address) return true;
if (gw.netmask < 32) {
@@ -149,7 +139,7 @@ module.exports = (srf, logger) => {
/* we have one or more carriers that match. Now we need to find one with a provisioned phone number */
const vc_sids = matches.map((m) => `'${m.voip_carrier_sid}'`).join(',');
const did = normalizeDID(req.calledNumber);
const sql = `SELECT * FROM phone_numbers WHERE number = '${did}' AND voip_carrier_sid IN (${vc_sids})`;
const sql = `SELECT * FROM phone_numbers WHERE number = ${did} AND voip_carrier_sid IN (${vc_sids})`;
logger.debug({matches, sql, did, vc_sids}, 'looking up DID');
const [r] = await pp.query(sql);
@@ -189,15 +179,6 @@ module.exports = (srf, logger) => {
}
}
/**
* The host part of the SIP URI is not a dot-decimal IP address,
* so this can be one of two things:
* (1) a sip realm value associate with an account, or
* (2) a carrier name for a carrier that we send outbound registrations to
*
* Let's look for case #1 first...
*/
/* get all the carriers and gateways for the account owning this sip realm */
const [gwAcc] = await pp.query(sqlSelectAllCarriersForAccountByRealm, uri.host);
const [gwSP] = gwAcc.length ? [[]] : await pp.query(sqlSelectAllCarriersForSPByRealm, uri.host);
@@ -215,41 +196,6 @@ module.exports = (srf, logger) => {
account: a[0]
};
}
/* no match, so let's look for case #2 */
try {
logger.info({
host: uri.host,
user: uri.user
}, 'sip realm is not associated with an account, checking carriers');
const [gw] = await pp.query(sqlSelectCarrierRequiringRegistration, [uri.host, uri.user]);
const matches = gw.filter(gatewayMatchesSourceAddress.bind(null, req.source_address));
if (1 === matches.length) {
// bingo
//TODO: this assumes the carrier is associate to an account, not an SP
//if the carrier is associated with an SP (which would mean we
//must see a dialed number in the To header, not the register username),
//then we need to look up the account based on the dialed number in the To header
const [a] = await pp.query(sqlAccountBySid, matches[0].account_sid);
if (0 === a.length) return failure;
logger.debug({matches}, `found registration carrier using ${uri.host} and ${uri.user}`);
return {
fromCarrier: true,
gateway: matches[0],
service_provider_sid: a[0].service_provider_sid,
account_sid: a[0].account_sid,
application_sid: matches[0].application_sid,
account: a[0]
};
}
else if (matches.length > 1) {
logger.warn({matches, source_address: req.source_address}, 'multiple gateways match source address');
}
} catch (err) {
logger.info({err, host: uri.host, user: uri.user}, 'Error looking up carrier by host and user');
}
return failure;
};
+2 -6
View File
@@ -2,7 +2,7 @@ const debug = require('debug')('jambonz:sbc-inbound');
const assert = require('assert');
const Emitter = require('events');
const parseUri = require('drachtio-srf').parseUri;
const {nudgeCallCounts, roundTripTime} = require('./utils');
const {nudgeCallCounts} = require('./utils');
const msProxyIps = process.env.MS_TEAMS_SIP_PROXY_IPS ?
process.env.MS_TEAMS_SIP_PROXY_IPS.split(',').map((i) => i.trim()) :
[];
@@ -135,8 +135,7 @@ module.exports = function(srf, logger) {
const identifyAccount = async(req, res, next) => {
try {
const {siprec, callId} = req.locals;
const {getSPForAccount, wasOriginatedFromCarrier, getApplicationForDidAndCarrier, stats} = req.srf.locals;
const startAt = process.hrtime();
const {getSPForAccount, wasOriginatedFromCarrier, getApplicationForDidAndCarrier} = req.srf.locals;
const {
fromCarrier,
gateway,
@@ -145,9 +144,6 @@ module.exports = function(srf, logger) {
service_provider_sid,
account
} = await wasOriginatedFromCarrier(req);
const rtt = roundTripTime(startAt);
stats.histogram('app.mysql.response_time', rtt, [
'query:wasOriginatedFromCarrier', 'app:sbc-inbound']);
/**
* calls come from 3 sources:
* (1) A carrier
+8 -28
View File
@@ -12,12 +12,10 @@ const getAppserver = (srf) => {
};
function makeRtpEngineOpts(req, srcIsUsingSrtp, dstIsUsingSrtp, teams = false) {
const rtpCopy = JSON.parse(JSON.stringify(rtpCharacteristics));
const srtpCopy = JSON.parse(JSON.stringify(srtpCharacteristics));
const from = req.getParsedHeader('from');
const srtpOpts = teams ? srtpCopy['teams'] : srtpCopy['default'];
const dstOpts = dstIsUsingSrtp ? srtpOpts : rtpCopy;
const srcOpts = srcIsUsingSrtp ? srtpOpts : rtpCopy;
const srtpOpts = teams ? srtpCharacteristics['teams'] : srtpCharacteristics['default'];
const dstOpts = dstIsUsingSrtp ? srtpOpts : rtpCharacteristics;
const srcOpts = srcIsUsingSrtp ? srtpOpts : rtpCharacteristics;
/* webrtc clients (e.g. sipjs) send DMTF via SIP INFO */
if ((srcIsUsingSrtp || dstIsUsingSrtp) && !teams) {
@@ -71,26 +69,6 @@ const systemHealth = async(redisClient, ping, getCount) => {
return getCount();
};
const doListen = (logger, app, port, resolve) => {
return app.listen(port, () => {
logger.info(`Health check server listening on http://localhost:${port}`);
resolve(app);
});
};
const handleErrors = (logger, app, resolve, reject, e) => {
if (e.code === 'EADDRINUSE' &&
process.env.HTTP_PORT_MAX &&
e.port < process.env.HTTP_PORT_MAX) {
logger.info(`Health check server failed to bind port on ${e.port}, will try next port`);
const server = doListen(logger, app, ++e.port, resolve);
server.on('error', handleErrors.bind(null, logger, app, resolve, reject));
return;
}
reject(e);
};
const createHealthCheckApp = (port, logger) => {
const express = require('express');
const app = express();
@@ -98,9 +76,11 @@ const createHealthCheckApp = (port, logger) => {
app.use(express.urlencoded({ extended: true }));
app.use(express.json());
return new Promise((resolve, reject) => {
const server = doListen(logger, app, port, resolve);
server.on('error', handleErrors.bind(null, logger, app, resolve, reject));
return new Promise((resolve) => {
app.listen(port, () => {
logger.info(`Health check server started at http://localhost:${port}`);
resolve(app);
});
});
};
+324 -1156
View File
File diff suppressed because it is too large Load Diff
+5 -5
View File
@@ -1,6 +1,6 @@
{
"name": "sbc-inbound",
"version": "v0.7.7",
"version": "v0.7.6",
"main": "app.js",
"engines": {
"node": ">= 12.0.0"
@@ -25,14 +25,14 @@
"jslint": "eslint app.js lib"
},
"dependencies": {
"@jambonz/db-helpers": "^0.7.3",
"@jambonz/realtimedb-helpers": "^0.5.7",
"@jambonz/db-helpers": "^0.6.19",
"@jambonz/http-authenticator": "^0.2.2",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/rtpengine-utils": "^0.3.11",
"@jambonz/realtimedb-helpers": "^0.4.29",
"@jambonz/rtpengine-utils": "^0.3.6",
"@jambonz/siprec-client-utils": "^0.1.4",
"@jambonz/stats-collector": "^0.1.6",
"@jambonz/time-series": "^0.2.5",
"@jambonz/time-series": "^0.2.4",
"aws-sdk": "^2.1152.0",
"bent": "^7.3.12",
"cidr-matcher": "^2.1.1",