Compare commits

..
8 changed files with 3307 additions and 4056 deletions
+36 -3
View File
@@ -14,6 +14,31 @@ assert.ok(process.env.DRACHTIO_SECRET, 'missing DRACHTIO_SECRET env var');
assert.ok(process.env.JAMBONES_TIME_SERIES_HOST, 'missing JAMBONES_TIME_SERIES_HOST env var');
assert.ok(process.env.JAMBONES_NETWORK_CIDR || process.env.K8S, 'missing JAMBONES_NETWORK_CIDR env var');
const JAMBONES_REDIS_SENTINELS = process.env.JAMBONES_REDIS_SENTINELS ? {
sentinels: process.env.JAMBONES_REDIS_SENTINELS.split(',').map((sentinel) => {
let host, port = 26379;
if (sentinel.includes(':')) {
const arr = sentinel.split(':');
host = arr[0];
port = parseInt(arr[1], 10);
} else {
host = sentinel;
}
return {host, port};
}),
name: process.env.JAMBONES_REDIS_SENTINEL_MASTER_NAME,
...(process.env.JAMBONES_REDIS_SENTINEL_PASSWORD && {
password: process.env.JAMBONES_REDIS_SENTINEL_PASSWORD
}),
...(process.env.JAMBONES_REDIS_SENTINEL_USERNAME && {
username: process.env.JAMBONES_REDIS_SENTINEL_USERNAME
}),
...(process.env.JAMBONES_REDIS_SENTINEL_SENTINAL_PASSWORD && {
sentinelPassword: process.env.JAMBONES_REDIS_SENTINEL_SENTINAL_PASSWORD
}),
} : null;
const Srf = require('drachtio-srf');
const srf = new Srf('sbc-inbound');
const opts = Object.assign({
@@ -42,6 +67,7 @@ const {LifeCycleEvents} = require('./lib/constants');
const setNameRtp = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-rtp`;
const rtpServers = [];
const setName = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-sip`;
const Registrar = require('@jambonz/mw-registrar');
const {
pool,
@@ -55,7 +81,8 @@ const {
lookupAccountBySid,
lookupAccountCapacitiesBySid,
queryCallLimits,
lookupClientByAccountAndUsername
lookupClientByAccountAndUsername,
lookupAppBySid
} = require('@jambonz/db-helpers')({
host: process.env.JAMBONES_MYSQL_HOST,
port: process.env.JAMBONES_MYSQL_PORT || 3306,
@@ -71,7 +98,11 @@ const {
addToSet,
removeFromSet,
incrKey,
decrKey} = require('@jambonz/realtimedb-helpers')({}, logger);
decrKey} = require('@jambonz/realtimedb-helpers')(JAMBONES_REDIS_SENTINELS || {
host: process.env.JAMBONES_REDIS_HOST,
port: process.env.JAMBONES_REDIS_PORT || 6379
}, logger);
const registrar = new Registrar(logger, redisClient);
const ngProtocol = process.env.JAMBONES_NG_PROTOCOL || 'udp';
const ngPort = process.env.RTPENGINE_PORT || ('udp' === ngProtocol ? 22222 : 8080);
@@ -92,6 +123,7 @@ srf.locals = {...srf.locals,
activeCallIds: new Map(),
getRtpEngine,
dbHelpers: {
registrar,
pool,
ping,
lookupAuthHook,
@@ -102,7 +134,8 @@ srf.locals = {...srf.locals,
lookupAccountBySipRealm,
lookupAccountCapacitiesBySid,
queryCallLimits,
lookupClientByAccountAndUsername
lookupClientByAccountAndUsername,
lookupAppBySid
},
realtimeDbHelpers: {
createSet,
+31 -110
View File
@@ -6,8 +6,7 @@ const {
SdpWantsSDES,
nudgeCallCounts,
roundTripTime,
parseConnectionIp,
isPrivateVoipNetwork
parseConnectionIp
} = require('./utils');
const {forwardInDialogRequests} = require('drachtio-fn-b2b-sugar');
@@ -15,7 +14,6 @@ const {parseUri, stringifyUri, SipError} = require('drachtio-srf');
const debug = require('debug')('jambonz:sbc-inbound');
const MS_TEAMS_USER_AGENT = 'Microsoft.PSTNHub.SIPProxy';
const MS_TEAMS_SIP_ENDPOINT = 'sip.pstnhub.microsoft.com';
const IMMUTABLE_HEADERS = ['via', 'from', 'to', 'call-id', 'cseq', 'max-forwards', 'content-length'];
/**
* this is to make sure the outgoing From has the number in the incoming From
@@ -67,7 +65,6 @@ class CallSession extends Emitter {
this.account_sid = req.locals.account_sid;
this.service_provider_sid = req.locals.service_provider_sid;
this.srsClients = [];
this.recordingNoAnswerTimeout = (process.env.JAMBONES_RECORDING_NO_ANSWER_TIMEOUT || 2) * 1000;
}
get isFromMSTeams() {
@@ -120,7 +117,6 @@ class CallSession extends Emitter {
offer,
answer,
del,
query,
blockMedia,
unblockMedia,
blockDTMF,
@@ -135,7 +131,6 @@ class CallSession extends Emitter {
this.offer = offer;
this.answer = answer;
this.del = del;
this.query = query;
this.blockMedia = blockMedia;
this.unblockMedia = unblockMedia;
this.blockDTMF = blockDTMF;
@@ -159,10 +154,7 @@ class CallSession extends Emitter {
const wantsSrtp = this.req.locals.possibleWebRtcClient = SdpWantsSrtp(sdp);
const wantsSDES = SdpWantsSDES(sdp);
this.rtpEngineOpts = makeRtpEngineOpts(this.req, wantsSrtp, false, this.isFromMSTeams || wantsSDES);
this.rtpEngineResource = {
destroy: this.del.bind(null, this.rtpEngineOpts.common),
query: this.query.bind(null, this.rtpEngineOpts.common),
};
this.rtpEngineResource = {destroy: this.del.bind(null, this.rtpEngineOpts.common)};
const obj = parseUri(this.req.uri);
let proxy, host, uri;
@@ -185,7 +177,7 @@ class CallSession extends Emitter {
...this.rtpEngineOpts.common,
...this.rtpEngineOpts.uac.mediaOpts,
'from-tag': this.rtpEngineOpts.uas.tag,
direction: [isPrivateVoipNetwork(this.req.source_address) ? 'private' : 'public', 'private'],
direction: ['public', 'private'],
sdp
};
const startAt = process.hrtime();
@@ -232,6 +224,12 @@ class CallSession extends Emitter {
if (this.req.locals.application_sid) {
Object.assign(headers, {'X-Application-Sid': this.req.locals.application_sid});
}
if (this.req.locals.queue_name) {
Object.assign(headers, {'X-Queue-Name': this.req.locals.queue_name});
}
if (this.req.locals.called_user) {
Object.assign(headers, {'X-Called-User': this.req.locals.called_user});
}
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});
@@ -256,8 +254,7 @@ class CallSession extends Emitter {
'-Max-Forwards',
'-Record-Route',
'-Session-Expires',
'-X-Application-Sid',
'-X-Authenticated-User'
'-X-Subspace-Forwarded-For',
],
proxyResponseHeaders: ['all', '-X-Trace-ID'],
localSdpB: spdOfferB,
@@ -336,22 +333,6 @@ class CallSession extends Emitter {
dlg.on('modify', this._onReinvite.bind(this, dlg));
}
_startRecordingNoAnswerTimer(res) {
this._clearRecordingNoAnswerTimer();
this.recordingNoAnswerTimer = setTimeout(() => {
this.logger.info('No response from SipRec server, return error to feature server');
this.isRecordingNoAnswerResponded = true;
res.send(400);
}, this.recordingNoAnswerTimeout);
}
_clearRecordingNoAnswerTimer() {
if (this.recordingNoAnswerTimer) {
clearTimeout(this.recordingNoAnswerTimer);
this.recordingNoAnswerTimer = null;
}
}
_stopRecording() {
if (this.srsClients.length) {
this.srsClients.forEach((c) => c.stop());
@@ -380,24 +361,12 @@ class CallSession extends Emitter {
this.uas = uas;
this.uac = uac;
[uas, uac].forEach((dlg) => {
dlg.on('destroy', async(bye) => {
dlg.on('destroy', async() => {
const other = dlg.other;
this.rtpEngineResource.destroy().catch((err) => {});
/* DH: need a better understanding of why query before delete is a good idea
this.rtpEngineResource.query()
.then((results) => {
this.logger.info({results}, 'rtpengine query results');
return this.rtpEngineResource.destroy();
})
.catch((err) => {});
*/
this.activeCallIds.delete(this.req.get('Call-ID'));
try {
const headers = {};
Object.keys(bye.headers).forEach((h) => {
if (!IMMUTABLE_HEADERS.includes(h)) headers[h] = bye.headers[h];
});
await other.destroy({headers});
await other.destroy();
} catch (err) {}
this.unsubscribeForDTMF();
@@ -693,7 +662,6 @@ Duration=${payload.duration} `
}
res.send(200, {body: response.sdp});
} catch (err) {
res.send(err.status || 500);
this.logger.error(err, 'Error handling reinvite');
}
}
@@ -724,12 +692,12 @@ Duration=${payload.duration} `
}
else if (reason.includes('CallRecording')) {
let succeeded = false;
const headers = contentType === 'application/json' && req.body ? JSON.parse(req.body) : {};
if (reason === 'startCallRecording') {
const from = this.req.getParsedHeader('From');
const to = this.req.getParsedHeader('To');
const aorFrom = from.uri;
const aorTo = to.uri;
const headers = contentType === 'application/json' && req.body ? JSON.parse(req.body) : {};
this.logger.info({to, from}, 'startCallRecording request for a call');
const srsUrl = req.get('X-Srs-Url');
@@ -773,97 +741,49 @@ Duration=${payload.duration} `
headers
}));
try {
this._startRecordingNoAnswerTimer(res);
await Promise.any(this.srsClients.map((c) => c.start()));
succeeded = true;
succeeded = (await Promise.all(
this.srsClients.map((c) => c.start())
)).every((r) => r);
} catch (err) {
this.logger.error({err}, 'Error starting SipRec call recording');
succeeded = false;
}
}
else if (reason === 'stopCallRecording') {
if (!this.srsClients.length || !this.srsClients.some((c) => c.activated)) {
if (!this.srsClients.length) {
res.send(400);
this.logger.info('discarding stopCallRecording request because we are not recording');
return;
}
try {
this._startRecordingNoAnswerTimer(res);
await Promise.any(this.srsClients.map((c) => {
if (c.activated) {
c.stop();
}
}));
succeeded = true;
succeeded = (await Promise.all(
this.srsClients.map((c) => c.stop())
)).every((r) => r);
} catch (err) {
this.logger.error({err}, 'Error stopping SipRec call recording');
succeeded = false;
}
this.srsClients = [];
}
else if (reason === 'pauseCallRecording') {
if (!this.srsClients.length || !this.srsClients.some((c) => c.activated && !c.paused)) {
if (!this.srsClients.length || this.srsClients.every((c) => c.paused)) {
this.logger.info('discarding invalid pauseCallRecording request');
res.send(400);
return;
}
try {
this._startRecordingNoAnswerTimer(res);
await Promise.any(this.srsClients.map((c) => {
if (c.activated && !c.paused) {
c.pause({headers});
}
}));
succeeded = true;
} catch (err) {
this.logger.error({err}, 'Error pausing SipRec call recording');
succeeded = false;
}
succeeded = (await Promise.all(
this.srsClients.map((c) => c.pause())
)).every((r) => r);
}
else if (reason === 'resumeCallRecording') {
if (!this.srsClients.length || !this.srsClients.some((c) => c.activated && c.paused)) {
if (!this.srsClients.length || !this.srsClients.every((c) => c.paused)) {
res.send(400);
this.logger.info('discarding invalid resumeCallRecording request');
return;
}
try {
this._startRecordingNoAnswerTimer(res);
await Promise.any(this.srsClients.map((c) => {
if (c.activated && c.paused) {
c.resume({headers});
}
}));
succeeded = true;
} catch (err) {
this.logger.error({err}, 'Error resuming SipRec call recording');
succeeded = false;
}
succeeded = (await Promise.all(
this.srsClients.map((c) => c.resume())
)).every((r) => r);
}
if (!this.isRecordingNoAnswerResponded) {
this._clearRecordingNoAnswerTimer();
res.send(succeeded ? 200 : 503);
}
} else if (reason.includes('Dtmf')) {
const arr = /Signal=\s*([0-9#*])/.exec(req.body);
if (!arr) {
this.logger.info({body: req.body}, '_onInfo: invalid INFO Dtmf');
throw new Error(`_onInfo: no dtmf in body for ${contentType}`);
}
const code = arr[1];
const arr2 = /Duration=\s*(\d+)/.exec(req.body);
const duration = arr2 ? arr2[1] : 250;
const dtmfOpts = {
...this.rtpEngineOpts.common,
'from-tag': this.rtpEngineOpts.uac.tag,
code,
duration
};
const response = await this.playDTMF(dtmfOpts);
if ('ok' !== response.result) {
this.logger.info({response}, `rtpengine play Dtmf failed with ${JSON.stringify(response)}`);
throw new Error('rtpengine failed: answer');
}
res.send(200);
res.send(succeeded ? 200 : 503);
}
}
else if (dlg.type === 'uas' && ['application/dtmf-relay', 'application/dtmf'].includes(contentType)) {
@@ -901,9 +821,10 @@ Duration=${payload.duration} `
}
}
else {
const immutableHdrs = ['via', 'from', 'to', 'call-id', 'cseq', 'max-forwards', 'content-length'];
const headers = {};
Object.keys(req.headers).forEach((h) => {
if (!IMMUTABLE_HEADERS.includes(h)) headers[h] = req.headers[h];
if (!immutableHdrs.includes(h)) headers[h] = req.headers[h];
});
const response = await dlg.other.request({method: 'INFO', headers, body: req.body});
const responseHeaders = {};
+4 -17
View File
@@ -142,7 +142,7 @@ module.exports = (srf, logger) => {
const failure = {fromCarrier: false};
const uri = parseUri(req.uri);
const isDotDecimal = /^(?:[0-9]{1,3}\.){3}[0-9]{1,3}$/.test(uri.host);
const did = normalizeDID(req.calledNumber) || 'anonymous';
if (!isDotDecimal) {
/**
* The host part of the SIP URI is not a dot-decimal IP address,
@@ -175,19 +175,12 @@ module.exports = (srf, logger) => {
.sort((a, b) => b.netmask - a.netmask);
const selected = gw.find(gatewayMatchesSourceAddress.bind(null, logger, req.source_address));
if (selected) {
const sql =
`SELECT application_sid FROM phone_numbers WHERE number = '${did}'
AND voip_carrier_sid = '${selected.voip_carrier_sid}'
AND account_sid = '${a[0].account_sid}'`;
logger.debug({selected, sql, did}, 'looking up DID');
const [r] = await pp.query(sql);
return {
fromCarrier: true,
gateway: selected,
service_provider_sid: a[0].service_provider_sid,
account_sid: a[0].account_sid,
application_sid: r[0]?.application_sid || selected.application_sid,
application_sid: selected.application_sid,
account: a[0]
};
}
@@ -214,19 +207,12 @@ module.exports = (srf, logger) => {
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}`);
const sql =
`SELECT application_sid FROM phone_numbers WHERE number = '${did}'
AND voip_carrier_sid = '${matches[0].voip_carrier_sid}'
AND account_sid = '${matches[0].account_sid}'`;
logger.debug({matches: matches[0], sql, did}, 'looking up DID');
const [r] = await pp.query(sql);
return {
fromCarrier: true,
gateway: matches[0],
service_provider_sid: a[0].service_provider_sid,
account_sid: a[0].account_sid,
application_sid: r[0]?.application_sid || matches[0].application_sid,
application_sid: matches[0].application_sid,
account: a[0]
};
}
@@ -278,6 +264,7 @@ module.exports = (srf, logger) => {
if (matches.length) {
/* we have one or more matches. Now check for one with a provisioned phone number matching the DID */
const vc_sids = matches.map((m) => `'${m.voip_carrier_sid}'`).join(',');
const did = normalizeDID(req.calledNumber) || 'anonymous';
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');
+25 -2
View File
@@ -28,7 +28,9 @@ module.exports = function(srf, logger) {
lookupAccountBySipRealm,
lookupAccountBySid,
lookupAccountCapacitiesBySid,
queryCallLimits
queryCallLimits,
lookupAppBySid,
registrar
} = srf.locals.dbHelpers;
const {stats, writeCdrs} = srf.locals;
@@ -196,11 +198,32 @@ module.exports = function(srf, logger) {
res.send(404);
return req.srf.endSession(req);
}
let deviceAppSid = null;
let called_user = null;
const queue_name = uri.user.startsWith('queue-') ? uri.user.match(/queue-(.*)/)[1] : null;
if (uri.user.startsWith('app-')) {
// Call from registered device to test application.
const appSid = uri.user.match(/app-(.*)/)[1];
const app = await lookupAppBySid(appSid);
if (app && app.account_sid === account.account_sid) {
deviceAppSid = app.application_sid;
}
} else if (!queue_name) {
// check if call to registered user
const {realm} = this.req.authorization.challengeResponse;
const calledAor = `${req.calledNumber}@${realm}`;
const reg = await registrar.query(calledAor);
if (reg) {
called_user = calledAor;
}
}
req.locals = {
service_provider_sid: account.service_provider_sid,
account_sid: account.account_sid,
account,
application_sid: account.device_calling_application_sid,
application_sid: deviceAppSid || account.device_calling_application_sid,
...(queue_name && ({queue_name})),
...(called_user && ({called_user})),
webhook_secret: account.webhook_secret,
realm: uri.host,
...(account.registration_hook && {
+3 -12
View File
@@ -22,8 +22,8 @@ function makeRtpEngineOpts(req, srcIsUsingSrtp, dstIsUsingSrtp, teams = false) {
const dstOpts = dstIsUsingSrtp ? srtpOpts : rtpCopy;
const srcOpts = srcIsUsingSrtp ? srtpOpts : rtpCopy;
/* Allow Feature server to inject DTMF to both leg except call from Teams*/
if (!teams) {
/* webrtc clients (e.g. sipjs) send DMTF via SIP INFO */
if ((srcIsUsingSrtp || dstIsUsingSrtp) && !teams) {
dstOpts.flags.push('inject DTMF');
srcOpts.flags.push('inject DTMF');
}
@@ -195,14 +195,6 @@ const isMSTeamsCIDR = (ip) => {
return matcher.contains(ip);
};
const isPrivateVoipNetwork = (ip) => {
if (process.env.PRIVATE_VOIP_NETWORK_CIDR) {
const matcher = new CIDRMatcher(process.env.PRIVATE_VOIP_NETWORK_CIDR.split(','));
return matcher.contains(ip);
}
return false;
};
module.exports = {
isWSS,
SdpWantsSrtp,
@@ -219,6 +211,5 @@ module.exports = {
nudgeCallCounts,
roundTripTime,
parseConnectionIp,
isMSTeamsCIDR,
isPrivateVoipNetwork
isMSTeamsCIDR
};
+3191 -3894
View File
File diff suppressed because it is too large Load Diff
+17 -16
View File
@@ -1,6 +1,6 @@
{
"name": "sbc-inbound",
"version": "0.9.0",
"version": "0.8.4",
"main": "app.js",
"engines": {
"node": ">= 12.0.0"
@@ -25,30 +25,31 @@
"jslint": "eslint app.js lib"
},
"dependencies": {
"@jambonz/db-helpers": "^0.9.3",
"@jambonz/db-helpers": "^0.9.1",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/realtimedb-helpers": "^0.8.8",
"@jambonz/rtpengine-utils": "^0.4.4",
"@jambonz/siprec-client-utils": "^0.2.7",
"@jambonz/realtimedb-helpers": "^0.8.6",
"@jambonz/rtpengine-utils": "^0.4.3",
"@jambonz/siprec-client-utils": "^0.2.6",
"@jambonz/stats-collector": "^0.1.9",
"@jambonz/time-series": "^0.2.8",
"@jambonz/digest-utils": "^0.0.5",
"@aws-sdk/client-sns": "^3.549.0",
"@aws-sdk/client-auto-scaling": "^3.549.0",
"@jambonz/time-series": "^0.2.5",
"@jambonz/digest-utils": "^0.0.3",
"@jambonz/mw-registrar": "^0.2.4",
"@aws-sdk/client-sns": "^3.360.0",
"@aws-sdk/client-auto-scaling": "^3.360.0",
"bent": "^7.3.12",
"cidr-matcher": "^2.1.1",
"debug": "^4.3.4",
"drachtio-fn-b2b-sugar": "0.1.0",
"drachtio-srf": "^4.5.31",
"express": "^4.19.2",
"pino": "^8.20.0",
"drachtio-fn-b2b-sugar": "0.0.12",
"drachtio-srf": "^4.5.21",
"express": "^4.18.1",
"pino": "^7.11.0",
"verify-aws-sns-signature": "^0.1.0",
"xml2js": "^0.6.2"
"xml2js": "^0.4.23"
},
"devDependencies": {
"eslint": "^7.32.0",
"eslint-plugin-promise": "^6.1.1",
"eslint-plugin-promise": "^4.3.1",
"nyc": "^15.1.0",
"tape": "^5.7.5"
"tape": "^4.15.1"
}
}
-2
View File
@@ -36,8 +36,6 @@
Subject: uac-pcap-carrier-success
Content-Type: application/sdp
Content-Length: [len]
X-Authenticated-User: xhoaluu@jambonz.org
X-Application-Sid: APP_ID_1
v=0
o=user1 53655765 2353687637 IN IP[local_ip_type] [local_ip]