Compare commits

...
19 Commits
Author SHA1 Message Date
Dave Horton b072524585 update deps, only subscribe for DTMF when client is possibly webrtc 2022-11-06 10:12:18 -05:00
Dave Horton 7fb9966d20 update to latest db-helpers 2022-11-05 10:42:14 -04:00
Dave Horton 42b3317206 update to db-helpers with caching fix 2022-11-01 20:40:53 -04:00
Dave Horton db66d98d0a write stat for db lookup 2022-11-01 18:57:31 -04:00
Dave Horton 26169517dc fix for dtmf handler when running multiple instances on same EC2 2022-11-01 15:40:37 -04:00
Dave Horton c78ec892b9 bugfix: running multiple instances under EC2 2022-10-31 14:45:20 -04:00
Dave Horton 6d237d43b3 update to db-helpers@0.7.0 with caching option 2022-10-31 11:44:53 -04:00
Dave Horton 44d298a4f3 Issue/58 (#59)
* feature: provide a range of health check ports, HTTP_PORT to HTTP_PORT_MAX, bind to first available

* feature: provide a range of aws sns ports, AWS_SNS_PORT to AWS_SNS_PORT_MAX, bind to first available

* AWS SNS port range 3010-3019
2022-10-30 13:08:07 -04:00
Dave Horton 030059a596 bugfix: transferring conference legs to different FS 2022-10-27 10:27:09 -04:00
Dave Horton 78b60525e2 Feature/support carrier domain in invite (#57)
* add support for incoming calls from carriers we register with, who then put their domain in the host part of incoming INVITE

* fix query to lookup registration carriers
2022-10-25 13:44:43 -04:00
Dave Horton dc9103cfa1 update deps 2022-10-23 15:27:18 -04:00
Dave Horton 4dd4247dc2 update package-lock.json 2022-10-23 12:20:28 -04:00
Dave Horton 055350903a update rtpengine-utils again 2022-10-23 11:25:06 -04:00
Dave Horton 7e83791ef0 update to latest rtpengine-utils 2022-10-22 22:17:02 -04:00
Dave Horton dbedd7419e block media going back to Five9 voicestream 2022-10-20 23:05:34 -04:00
Markus FrindtandMarkus Frindt c4d07b517e [snyk] fix vulnerability (#56)
Co-authored-by: Markus Frindt <m.frindt@cognigy.com>
2022-10-20 21:35:55 -04:00
Dave Horton 2693077871 include X-Account-Sid header when moving call between FS 2022-10-15 10:58:27 -04:00
Dave Horton ad0912d302 bugfix for release media (#55) 2022-10-14 12:48:22 -04:00
Dave Horton e715433534 bump version 2022-10-13 16:00:41 -04:00
9 changed files with 1321 additions and 353 deletions
+1 -1
View File
@@ -1,4 +1,4 @@
FROM --platform=linux/amd64 node:18.8.0-alpine as base
FROM --platform=linux/amd64 node:18.9.0-alpine3.16 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
});
+28 -6
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 || 3001;
const PORT = process.env.AWS_SNS_PORT || 3010;
const {LifeCycleEvents} = require('./constants');
const express = require('express');
const app = express();
@@ -21,6 +21,26 @@ 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 {
@@ -45,6 +65,7 @@ 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':
@@ -80,14 +101,12 @@ class SnsNotifier extends Emitter {
async init() {
try {
this.logger.info('SnsNotifier: retrieving instance data');
this.logger.debug('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,
snsEndpoint: this.snsEndpoint
publicIp: this.publicIp
}, 'retrieved AWS instance data');
// start listening
@@ -99,7 +118,10 @@ class SnsNotifier extends Emitter {
this.logger.error(err, 'burped error');
res.status(err.status || 500).json({msg: err.message});
});
app.listen(PORT);
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));
});
} catch (err) {
this.logger.error({err}, 'Error retrieving AWS instance metadata');
+55 -17
View File
@@ -80,6 +80,28 @@ 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');
@@ -129,7 +151,8 @@ class CallSession extends Emitter {
}
this.logger.debug(`using feature server ${featureServer}`);
this.rtpEngineOpts = makeRtpEngineOpts(this.req, SdpWantsSrtp(sdp), false, this.isFromMSTeams);
const wantsSrtp = this.req.locals.possibleWebRtcClient = SdpWantsSrtp(sdp);
this.rtpEngineOpts = makeRtpEngineOpts(this.req, wantsSrtp, false, this.isFromMSTeams);
this.rtpEngineResource = {destroy: this.del.bind(null, this.rtpEngineOpts.common)};
const obj = parseUri(this.req.uri);
let proxy, host, uri;
@@ -246,6 +269,17 @@ class CallSession extends Emitter {
this.logger.error(`rtpengine answer failed with ${JSON.stringify(response)}`);
throw new Error('rtpengine failed: answer');
}
/* special case: Five9 Voicestream calls do not advertise a:sendonly, though they should */
if (this.isFive9VoiceStream) {
const opts = {
...this.rtpEngineOpts.common,
'from-tag':this.rtpEngineOpts.uac.tag
};
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');
}
return response.sdp;
}
});
@@ -275,8 +309,7 @@ class CallSession extends Emitter {
_setDlgHandlers(dlg) {
const {callId} = dlg.sip;
this.activeCallIds.set(callId, this);
this.subscribeDTMF(this.logger, callId, this.rtpEngineOpts.uas.tag,
this._onDTMF.bind(this));
if (this.isPossibleWebRtcClient) this.subscribeForDTMF(dlg);
dlg.on('destroy', () => {
debug('call ended with normal termination');
this.logger.info('call ended with normal termination');
@@ -324,7 +357,8 @@ class CallSession extends Emitter {
try {
await other.destroy();
} catch (err) {}
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
this.unsubscribeForDTMF();
const trackingOn = process.env.JAMBONES_TRACK_ACCOUNT_CALLS ||
process.env.JAMBONES_TRACK_SP_CALLS ||
@@ -380,8 +414,7 @@ class CallSession extends Emitter {
});
});
this.subscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag,
this._onDTMF.bind(this, uac));
if (this.isPossibleWebRtcClient) this.subscribeForDTMF(uac);
uas.on('modify', this._onReinvite.bind(this, uas));
uac.on('modify', this._onReinvite.bind(this, uac));
@@ -475,7 +508,7 @@ Duration=${payload.duration} `
});
}
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
this.unsubscribeForDTMF();
const uas = await this.srf.createUAS(req, res, {
localSdp: response.sdp,
@@ -508,11 +541,16 @@ 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,
@@ -531,11 +569,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 (reason && dlg.type === 'uac' && ['release-media', 'anchor-media'].includes(reason) &&
!this.callerIsUsingSrtp) {
if (isReleasingMedia && !this.callerIsUsingSrtp) {
this.logger.info({response}, `got a reinvite from FS to ${reason}`);
sdp = dlg.other.remote.sdp;
answerMedia.flags = ['asymmetric', 'port latching'];
if (!answerMedia.flags.includes('asymmetric')) answerMedia.flags.push('asymmetric');
answerMedia.flags = answerMedia.flags.filter((f) => f !== 'media handover');
this._mediaReleased = 'release-media' === reason;
}
else {
@@ -776,15 +814,15 @@ Duration=${payload.duration} `
res.send(202);
// invite to new fs
const headers = {};
if (req.has('X-Retain-Call-Sid')) {
Object.assign(headers, {'X-Retain-Call-Sid': req.get('X-Retain-Call-Sid')});
}
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 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._onFeatureServerReinvite.bind(this, uac));
uac.on('modify', this._onReinvite.bind(this, uac));
uac.on('refer', this._onFeatureServerTransfer.bind(this, uac));
uac.on('destroy', () => {
this.logger.info('call ended with normal termination');
@@ -802,7 +840,7 @@ Duration=${payload.duration} `
const response = await this.answer(opts);
if ('ok' !== response.result) {
res.send(488);
throw new Error(`_onFeatureServerReinvite: rtpengine failed: ${JSON.stringify(response)}`);
throw new Error(`_onFeatureServerTransfer: rtpengine failed: ${JSON.stringify(response)}`);
}
this.logger.info('successfully moved call to new feature server');
} catch (err) {
@@ -863,7 +901,7 @@ Duration=${payload.duration} `
// successfully connected
this.logger.info('successfully connected new call leg for REFER');
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
this.unsubscribeForDTMF();
this.referInvite = null;
sendNotify(this.uas, '200 OK');
this.uas.destroy();
+54
View File
@@ -57,6 +57,16 @@ 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) {
@@ -179,6 +189,15 @@ 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);
@@ -196,6 +215,41 @@ 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;
};
+6 -2
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} = require('./utils');
const {nudgeCallCounts, roundTripTime} = require('./utils');
const msProxyIps = process.env.MS_TEAMS_SIP_PROXY_IPS ?
process.env.MS_TEAMS_SIP_PROXY_IPS.split(',').map((i) => i.trim()) :
[];
@@ -135,7 +135,8 @@ module.exports = function(srf, logger) {
const identifyAccount = async(req, res, next) => {
try {
const {siprec, callId} = req.locals;
const {getSPForAccount, wasOriginatedFromCarrier, getApplicationForDidAndCarrier} = req.srf.locals;
const {getSPForAccount, wasOriginatedFromCarrier, getApplicationForDidAndCarrier, stats} = req.srf.locals;
const startAt = process.hrtime();
const {
fromCarrier,
gateway,
@@ -144,6 +145,9 @@ 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
+23 -5
View File
@@ -71,6 +71,26 @@ 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();
@@ -78,11 +98,9 @@ const createHealthCheckApp = (port, logger) => {
app.use(express.urlencoded({ extended: true }));
app.use(express.json());
return new Promise((resolve) => {
app.listen(port, () => {
logger.info(`Health check server started at http://localhost:${port}`);
resolve(app);
});
return new Promise((resolve, reject) => {
const server = doListen(logger, app, port, resolve);
server.on('error', handleErrors.bind(null, logger, app, resolve, reject));
});
};
+1149 -317
View File
File diff suppressed because it is too large Load Diff
+4 -4
View File
@@ -1,6 +1,6 @@
{
"name": "sbc-inbound",
"version": "v0.7.6",
"version": "v0.7.7",
"main": "app.js",
"engines": {
"node": ">= 12.0.0"
@@ -25,11 +25,11 @@
"jslint": "eslint app.js lib"
},
"dependencies": {
"@jambonz/db-helpers": "^0.6.19",
"@jambonz/db-helpers": "^0.7.3",
"@jambonz/realtimedb-helpers": "^0.5.7",
"@jambonz/http-authenticator": "^0.2.2",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/realtimedb-helpers": "^0.4.29",
"@jambonz/rtpengine-utils": "^0.3.6",
"@jambonz/rtpengine-utils": "^0.3.11",
"@jambonz/siprec-client-utils": "^0.1.4",
"@jambonz/stats-collector": "^0.1.6",
"@jambonz/time-series": "^0.2.5",