Compare commits

...
34 Commits
Author SHA1 Message Date
Dave Horton 43ed2bafe1 strip X-Preferred-From-User and X-Preferred-From-Host from outgoing call 2022-10-26 09:35:11 -04:00
Dave Horton 5ee041bdee feature: specify user or host part of From on outdial 2022-10-23 15:27:03 -04:00
Dave Horton 5bbe6a8752 update package-lock.json 2022-10-23 12:20:12 -04:00
Dave Horton 2e5d609bab update rtpengine-utils again 2022-10-23 11:27:34 -04:00
Dave Horton 8f598bf7b0 update to latest rtpengine-utils 2022-10-22 22:26:02 -04:00
Markus FrindtandMarkus Frindt c7d717b3ee [snyk] fix vulnerability (#55)
Co-authored-by: Markus Frindt <m.frindt@cognigy.com>
2022-10-20 21:36:07 -04:00
Dave Horton e9209b37ca include X-Account-Sid header when moving call between FS 2022-10-15 10:59:29 -04:00
Dave Horton 3b98dc6ec2 Bugfix/release media investigation (#54)
* added logging

* when releasing media, include asymettric flag on offer

* update package-lock.json

* remove media handover flag when releasing media
2022-10-14 12:46:14 -04:00
Dave Horton 1c84dd799c update time-series 2022-10-10 09:16:42 +01:00
Dave Horton 4338ae9411 bugfix: inject DMTF flag was inserted over and over 2022-10-07 11:52:36 +01:00
Dave Horton e17e7dbddd add support for connecting to rtpengine via ws 2022-09-27 09:43:05 +01:00
Dave Horton 254479e289 Feature/app call count tracking (#52)
* add call count tracking at the app level (optional)

* update call_counts_app schema

* update rtpengine-utils
2022-09-22 23:44:58 +02:00
Dave Horton a10a311dcb only track service provider calls if JAMBONES_TRACK_SP_CALLS is set 2022-09-20 16:42:22 +02:00
Dave Horton 806cb89c37 include application_sid in cdr 2022-09-20 14:00:33 +02:00
Dave Horton fffa2748d1 Feature/sp limits (#51)
* add account and service provider call limits

* add custom headers when rejecting calls due to max calls limit
2022-09-20 13:13:24 +02:00
Dave Horton 76625c7596 Feature/recent calls enhancement with sp (#50)
* write cdrs with service_provider_sid

* write call counts at the SP level
2022-09-16 13:45:33 +02:00
Paulo Tellesandp.souza 3b0f7ff6eb change node image (#48)
Co-authored-by: p.souza <p.souza@cognigy.com>
2022-09-07 13:20:13 +02:00
Dave Horton 33c75acc9e bump version to 0.7.6 2022-08-26 20:09:43 +02:00
xquanluu bae9ef7638 feat: update time-series 0.11.12 (#47) 2022-08-19 16:21:11 +02:00
Dave Horton 4ed4b38301 update time-series 2022-08-19 09:57:48 +02:00
Dave Horton 6b6f89264f minor logging change 2022-08-17 16:11:10 +02:00
Dave Horton f0e0fba2f1 make rtpengine transcode if non-preferred codec is selected by far end 2022-08-17 14:14:45 +02:00
Dave Horton 290723f234 initial changes for sip info dtmf (#46) 2022-08-11 15:13:07 +02:00
Dave Horton 5a14aa807a update to latest @jambonz/siprec-utils with fix for sdp version 2022-08-08 16:27:59 +02:00
Dave Horton 2505a36db6 fix 2022-08-05 10:11:21 +01:00
Dave Horton 90818206f5 Dockerfile: update base image 2022-08-05 10:10:24 +01:00
Dave Horton 1de4db6ebc Dockerfile: update base image 2022-07-28 12:58:23 +01:00
Dave Horton 7d2125788f when releasing media, use asymetric flag so that rtpengine does react to a spurious final packet from freeswitch by incorrectly sending rtp there 2022-07-26 12:25:58 +01:00
Snyk bot 3b83c1bda8 fix: upgrade drachtio-srf from 4.5.0 to 4.5.1 (#42)
Snyk has created this PR to upgrade drachtio-srf from 4.5.0 to 4.5.1.

See this package in npm:
https://www.npmjs.com/package/drachtio-srf

See this project in Snyk:
https://app.snyk.io/org/davehorton/project/282b1881-3e13-4fad-85bd-3b1c662c2134?utm_source=github&utm_medium=referral&page=upgrade-pr
2022-07-13 09:16:37 +02:00
Paulo Tellesandp.souza f352bf885c improve dockerfile to fix snyk security issues (#41)
Co-authored-by: p.souza <p.souza@cognigy.com>
2022-07-07 15:19:14 +02:00
Dave Horton 08a4f5defb update Dockerfile 2022-06-23 09:57:38 -04:00
Dave Horton 87df38110b update to azure 1.22.0 2022-06-11 16:23:08 -04:00
Dave Horton 06e370fa59 update deps 2022-06-11 11:55:21 -04:00
Dave Horton e4ed2cea26 healthcheck improvements (#38) 2022-04-12 15:45:33 -04:00
9 changed files with 2680 additions and 2669 deletions
+19 -6
View File
@@ -1,10 +1,23 @@
FROM node:17-slim
FROM --platform=linux/amd64 node:18.9.0-alpine3.16 as base
RUN apk --update --no-cache add --virtual .builds-deps build-base python3
WORKDIR /opt/app/
COPY package.json ./
RUN npm install
RUN npm prune
COPY . /opt/app
FROM base as build
COPY package.json package-lock.json ./
RUN npm ci
COPY . .
FROM base
COPY --from=build /opt/app /opt/app/
ARG NODE_ENV
ENV NODE_ENV $NODE_ENV
CMD [ "npm", "start" ]
CMD [ "node", "app.js" ]
+43 -11
View File
@@ -11,12 +11,15 @@ assert.ok(process.env.JAMBONES_NETWORK_CIDR || process.env.K8S, 'missing JAMBONE
const Srf = require('drachtio-srf');
const srf = new Srf('sbc-outbound');
const CIDRMatcher = require('cidr-matcher');
const {pingMsTeamsGateways, equalsIgnoreOrder} = require('./lib/utils');
const {equalsIgnoreOrder, pingMsTeamsGateways, createHealthCheckApp, systemHealth} = require('./lib/utils');
const opts = Object.assign({
timestamp: () => {return `, "time": "${new Date().toISOString()}"`;}
}, {level: process.env.JAMBONES_LOGLEVEL || 'info'});
const logger = require('pino')(opts);
const {
writeCallCount,
writeCallCountSP,
writeCallCountApp,
writeCdrs,
queryCdrs,
writeAlerts,
@@ -32,13 +35,15 @@ const CallSession = require('./lib/call-session');
const setNameRtp = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-rtp`;
const rtpServers = [];
const {
ping,
performLcr,
lookupAllTeamsFQDNs,
lookupAccountBySipRealm,
lookupAccountBySid,
lookupAccountCapacitiesBySid,
lookupSipGatewaysByCarrier,
lookupCarrierBySid
lookupCarrierBySid,
queryCallLimits
} = require('@jambonz/db-helpers')({
host: process.env.JAMBONES_MYSQL_HOST,
user: process.env.JAMBONES_MYSQL_USER,
@@ -47,6 +52,7 @@ const {
connectionLimit: process.env.JAMBONES_MYSQL_CONNECTION_LIMIT || 10
}, logger);
const {
client: redisClient,
createHash,
retrieveHash,
incrKey,
@@ -62,19 +68,24 @@ const activeCallIds = new Map();
srf.locals = {...srf.locals,
stats,
writeCallCount,
writeCallCountSP,
writeCallCountApp,
writeCdrs,
writeAlerts,
AlertType,
queryCdrs,
activeCallIds,
dbHelpers: {
ping,
performLcr,
lookupAllTeamsFQDNs,
lookupAccountBySipRealm,
lookupAccountBySid,
lookupAccountCapacitiesBySid,
lookupSipGatewaysByCarrier,
lookupCarrierBySid
lookupCarrierBySid,
queryCallLimits
},
realtimeDbHelpers: {
createHash,
@@ -88,10 +99,12 @@ const {initLocals, checkLimits, route} = require('./lib/middleware')(srf, logger
host: process.env.JAMBONES_REDIS_HOST,
port: process.env.JAMBONES_REDIS_PORT || 6379
});
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 || 22225,
protocol: 'udp'
protocol: ngProtocol
});
srf.locals.getRtpEngine = getRtpEngine;
@@ -137,13 +150,32 @@ srf.invite((req, res) => {
session.connect();
});
if (process.env.K8S) {
if (process.env.K8S || process.env.HTTP_PORT) {
const PORT = process.env.HTTP_PORT || 3000;
const getCount = () => activeCallIds.size;
const healthCheck = require('@jambonz/http-health-check');
healthCheck({port: PORT, logger, path: '/', fn: getCount});
}
const getCount = () => srf.locals.activeCallIds.size;
createHealthCheckApp(PORT, logger)
.then((app) => {
healthCheck({
app,
logger,
path: '/',
fn: getCount
});
healthCheck({
app,
logger,
path: '/system-health',
fn: systemHealth.bind(null, redisClient, ping, getCount)
});
return;
})
.catch((err) => {
logger.error({err}, 'Error creating health check server');
});
}
if ('test' !== process.env.NODE_ENV) {
/* update call stats periodically */
setInterval(() => {
@@ -164,7 +196,7 @@ const lookupRtpServiceEndpoints = (lookup, serviceName) => {
rtpServers.length = 0;
Array.prototype.push.apply(rtpServers, addrs);
logger.info({rtpServers}, 'rtpserver endpoints have been updated');
setRtpEngines(rtpServers.map((a) => `${a}:${process.env.RTPENGINE_PORT || 22222}`));
setRtpEngines(rtpServers.map((a) => `${a}:${ngPort}`));
}
});
};
@@ -191,7 +223,7 @@ else {
logger.debug({newArray, rtpServers}, 'getActiveRtpServers');
if (!equalsIgnoreOrder(newArray, rtpServers)) {
logger.info({newArray}, 'resetting active rtpengines');
setRtpEngines(newArray.map((a) => `${a}:${process.env.RTPENGINE_PORT || 22222}`));
setRtpEngines(newArray.map((a) => `${a}:${ngPort}`));
rtpServers.length = 0;
Array.prototype.push.apply(rtpServers, newArray);
}
+196 -25
View File
@@ -1,5 +1,7 @@
const Emitter = require('events');
const {makeRtpEngineOpts, makeCallCountKey} = require('./utils');
const sdpTransform = require('sdp-transform');
const SrsClient = require('@jambonz/siprec-client-utils');
const {makeRtpEngineOpts, nudgeCallCounts} = require('./utils');
const {forwardInDialogRequests} = require('drachtio-fn-b2b-sugar');
const {SipError, stringifyUri, parseUri} = require('drachtio-srf');
const debug = require('debug')('jambonz:sbc-outbound');
@@ -11,10 +13,17 @@ const makeInviteInProgressKey = (callid) => `sbc-out-iip${callid}`;
*/
const createBLegFromHeader = (req, teams) => {
const from = req.getParsedHeader('From');
const host = teams ? req.get('X-MS-Teams-Tenant-FQDN') : 'localhost';
const uri = parseUri(from.uri);
if (uri && uri.user) return `sip:${uri.user}@${host}`;
return `sip:anonymous@${host}`;
let user = uri.user || 'anonymous';
let host = 'localhost';
if (teams) {
host = req.get('X-MS-Teams-Tenant-FQDN');
}
else if (req.has('X-Preferred-From-User') || req.has('X-Preferred-From-Host')) {
user = req.get('X-Preferred-From-User') || user;
host = req.get('X-Preferred-From-Host') || host;
}
return `sip:${user}@${host}`;
};
const createBLegToHeader = (req, teams) => {
const to = req.getParsedHeader('To');
@@ -48,6 +57,15 @@ const initCdr = (srf, req) => {
};
};
const updateRtpEngineFlags = (sdp, opts) => {
try {
const parsed = sdpTransform.parse(sdp);
const codec = parsed.media[0].rtp[0].codec;
if (['PCMU', 'PCMA'].includes(codec)) opts.flags.push(`codec-accept-${codec}`);
} catch (err) {}
return opts;
};
class CallSession extends Emitter {
constructor(logger, req, res) {
super();
@@ -60,24 +78,32 @@ class CallSession extends Emitter {
this.activeCallIds = this.srf.locals.activeCallIds;
this.writeCdrs = this.srf.locals.writeCdrs;
this.incrKey = req.srf.locals.realtimeDbHelpers.incrKey;
this.decrKey = req.srf.locals.realtimeDbHelpers.decrKey;
this.callCountKey = makeCallCountKey(req.locals.account_sid);
const {performLcr, lookupCarrierBySid, lookupSipGatewaysByCarrier} = this.srf.locals.dbHelpers;
this.performLcr = performLcr;
this.lookupCarrierBySid = lookupCarrierBySid;
this.lookupSipGatewaysByCarrier = lookupSipGatewaysByCarrier;
this._mediaReleased = false;
}
get account_sid() {
return this.req.locals.account_sid;
}
get application_sid() {
return this.req.locals.application_sid;
}
get privateSipAddress() {
return this.srf.locals.privateSipAddress;
}
get isMediaReleased() {
return this._mediaReleased;
}
async connect() {
const teams = this.teams = this.req.locals.target === 'teams';
const engine = this.srf.locals.getRtpEngine();
@@ -94,8 +120,12 @@ class CallSession extends Emitter {
unblockMedia,
blockDTMF,
unblockDTMF,
playDTMF,
subscribeDTMF,
unsubscribeDTMF
unsubscribeDTMF,
subscribeRequest,
subscribeAnswer,
unsubscribe
} = engine;
const {createHash, retrieveHash} = this.srf.locals.realtimeDbHelpers;
this.offer = offer;
@@ -105,8 +135,12 @@ class CallSession extends Emitter {
this.unblockMedia = unblockMedia;
this.blockDTMF = blockDTMF;
this.unblockDTMF = unblockDTMF;
this.playDTMF = playDTMF;
this.subscribeDTMF = subscribeDTMF;
this.unsubscribeDTMF = unsubscribeDTMF;
this.subscribeRequest = subscribeRequest;
this.subscribeAnswer = subscribeAnswer;
this.unsubscribe = unsubscribe;
this.rtpEngineOpts = makeRtpEngineOpts(this.req, false, this.useWss || teams, teams);
this.rtpEngineResource = {destroy: this.del.bind(null, this.rtpEngineOpts.common)};
@@ -221,13 +255,13 @@ class CallSession extends Emitter {
}
// rtpengine 'offer'
const opts = {
const opts = updateRtpEngineFlags(this.req.body, {
...this.rtpEngineOpts.common,
...this.rtpEngineOpts.uac.mediaOpts,
'from-tag': this.rtpEngineOpts.uas.tag,
direction: ['private', 'public'],
sdp: this.req.body
};
});
const response = await this.offer(opts);
debug(`response from rtpengine to offer ${JSON.stringify(response)}`);
this.logger.debug({offer: opts, response}, 'initial offer to rtpengine');
@@ -288,13 +322,14 @@ class CallSession extends Emitter {
'-X-MS-Teams-FQDN',
'-X-MS-Teams-Tenant-FQDN',
'-X-Trace-ID',
'X-CID',
'-Allow',
'-Session-Expires',
'-X-Requested-Carrier-Sid',
'-X-Jambonz-Routing',
'-X-Jambonz-FS-UUID',
'Min-SE'
'-X-Preferred-From-User',
'X-Preferred-From-Host',
'-X-Jambonz-FS-UUID',
],
proxyResponseHeaders: [
'all',
@@ -336,7 +371,9 @@ class CallSession extends Emitter {
if (!this.req.locals.account.disable_cdrs) {
this.req.locals.cdr = {
...initCdr(this.req.srf, inv),
service_provider_sid: this.req.locals.service_provider_sid,
account_sid: this.req.locals.account_sid,
...(this.req.locals.application_sid && {application_sid: this.req.locals.application_sid}),
trunk
};
}
@@ -430,13 +467,19 @@ class CallSession extends Emitter {
await other.destroy();
} catch (err) {}
this.decrKey(this.callCountKey)
.then((count) => {
this.logger.debug(`after hangup there are ${count} active calls for this account`);
debug(`after hangup there are ${count} active calls for this account`);
return;
})
.catch((err) => this.logger.error({err}, 'Error decrementing call count'));
const trackingOn = process.env.JAMBONES_TRACK_ACCOUNT_CALLS ||
process.env.JAMBONES_TRACK_SP_CALLS ||
process.env.JAMBONES_TRACK_APP_CALLS;
if (process.env.JAMBONES_HOSTING || trackingOn) {
const {writeCallCount, writeCallCountSP, writeCallCountApp} = this.req.srf.locals;
await nudgeCallCounts(this.logger, {
service_provider_sid: this.service_provider_sid,
account_sid: this.account_sid,
application_sid: this.application_sid
}, this.decrKey, {writeCallCountSP, writeCallCount, writeCallCountApp})
.catch((err) => this.logger.error(err, 'Error decrementing call counts'));
}
/* write cdr for connected call */
if (this.req.locals.cdr) {
@@ -525,11 +568,17 @@ Duration=${payload.duration} `
async _onReinvite(dlg, req, res) {
try {
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' ? ['private', 'public'] : ['public', 'private'];
if (isReleasingMedia) {
if (!offerMedia.flags.includes('port latching')) offerMedia.flags.push('port latching');
if (!offerMedia.flags.includes('asymmetric')) offerMedia.flags.push('asymmetric');
offerMedia.flags = offerMedia.flags.filter((f) => f !== 'media handover');
}
let opts = {
...this.rtpEngineOpts.common,
...offerMedia,
@@ -540,18 +589,22 @@ Duration=${payload.duration} `
};
if (reason && opts.flags && !opts.flags.includes('reset')) opts.flags.push('reset');
let response = await this.offer(opts);
if ('ok' !== response.result) {
res.send(488);
throw new Error(`_onReinvite: rtpengine failed: offer: ${JSON.stringify(response)}`);
}
this.logger.debug({opts, response}, 'CallSession:_onReinvite: (offer)');
/* if this is a re-invite from the FS to change media anchoring, avoid sending the reinvite out */
let sdp;
if (reason && dlg.type === 'uas' && ['release-media', 'anchor-media'].includes(reason)) {
if (isReleasingMedia) {
this.logger.info(`got a reinvite from FS to ${reason}`);
sdp = dlg.other.remote.sdp;
if (!answerMedia.flags.includes('port latching')) answerMedia.flags.push('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 {
sdp = await dlg.other.modify(response.sdp);
@@ -569,7 +622,7 @@ Duration=${payload.duration} `
res.send(488);
throw new Error(`_onReinvite: rtpengine failed: ${JSON.stringify(response)}`);
}
this.logger.info({sdp: response.sdp}, 'CallSession:_onReinvite: sending back upstream');
this.logger.debug({opts, sdp: response.sdp}, 'CallSession:_onReinvite: (answer) sending back upstream');
res.send(200, {body: response.sdp});
} catch (err) {
this.logger.error(err, 'Error handling reinvite');
@@ -578,6 +631,7 @@ Duration=${payload.duration} `
async _onInfo(dlg, req, res) {
try {
const contentType = req.get('Content-Type');
if (dlg.type === 'uas' && req.has('X-Reason')) {
const toTag = this.rtpEngineOpts.uac.tag;
const reason = req.get('X-Reason');
@@ -596,6 +650,123 @@ Duration=${payload.duration} `
const response = Promise.all([this.unblockMedia(opts), this.unblockDTMF(opts)]);
this.logger.info({response}, `_onInfo: response to rtpengine command for ${reason}`);
}
else if (reason.includes('CallRecording')) {
let succeeded = false;
if (reason === 'startCallRecording') {
const from = this.req.getParsedHeader('From');
const to = this.req.getParsedHeader('To');
const aorFrom = from.uri;
const aorTo = to.uri;
this.logger.info({to, from}, 'startCallRecording request for an outbound call');
const srsUrl = req.get('X-Srs-Url');
const srsRecordingId = req.get('X-Srs-Recording-ID');
const callSid = req.get('X-Call-Sid');
const accountSid = req.get('X-Account-Sid');
const applicationSid = req.get('X-Application-Sid');
if (this.srsClient) {
res.send(400);
this.logger.info('discarding duplicate startCallRecording request for a call');
return;
}
if (!srsUrl) {
this.logger.info('startCallRecording request is missing X-Srs-Url header');
res.send(400);
return;
}
this.srsClient = new SrsClient(this.logger, {
srf: dlg.srf,
direction: 'outbound',
originalInvite: this.req,
callingNumber: this.req.callingNumber,
calledNumber: this.req.calledNumber,
srsUrl,
srsRecordingId,
callSid,
accountSid,
applicationSid,
rtpEngineOpts: this.rtpEngineOpts,
toTag,
aorFrom,
aorTo,
subscribeRequest: this.subscribeRequest,
subscribeAnswer: this.subscribeAnswer,
del: this.del,
blockMedia: this.blockMedia,
unblockMedia: this.unblockMedia,
unsubscribe: this.unsubscribe
});
try {
succeeded = await this.srsClient.start();
} catch (err) {
this.logger.error({err}, 'Error starting SipRec call recording');
}
}
else if (reason === 'stopCallRecording') {
if (!this.srsClient) {
res.send(400);
this.logger.info('discarding stopCallRecording request because we are not recording');
return;
}
try {
succeeded = await this.srsClient.stop();
} catch (err) {
this.logger.error({err}, 'Error stopping SipRec call recording');
}
this.srsClient = null;
}
else if (reason === 'pauseCallRecording') {
if (!this.srsClient || this.srsClient.paused) {
this.logger.info('discarding invalid pauseCallRecording request');
res.send(400);
return;
}
succeeded = await this.srsClient.pause();
}
else if (reason === 'resumeCallRecording') {
if (!this.srsClient || !this.srsClient.paused) {
res.send(400);
this.logger.info('discarding invalid resumeCallRecording request');
return;
}
succeeded = await this.srsClient.resume();
}
res.send(succeeded ? 200 : 503);
}
}
else if (dlg.type === 'uac' && ['application/dtmf-relay', 'application/dtmf'].includes(contentType)) {
const arr = /Signal=\s*([1-9#*])/.exec(req.body);
if (!arr) {
this.logger.info({body: req.body}, '_onInfo: invalid INFO dtmf request');
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;
if (this.isMediaReleased) {
/* just relay on to the feature server */
this.logger.info({code, duration}, 'got SIP INFO DTMF from caller, relaying to feature server');
this._onDTMF(dlg.other, {event: code, duration})
.catch((err) => this.logger.info({err}, 'Error relaying DTMF to feature server'));
res.send(200);
}
else {
/* else convert SIP INFO to RFC 2833 telephony events */
this.logger.info({code, duration}, 'got SIP INFO DTMF from caller, converting to RFC 2833');
const opts = {
...this.rtpEngineOpts.common,
'from-tag': this.rtpEngineOpts.uac.tag,
code,
duration
};
const response = await this.playDTMF(opts);
if ('ok' !== response.result) {
this.logger.info({response}, `rtpengine playDTMF failed with ${JSON.stringify(response)}`);
throw new Error('rtpengine failed: answer');
}
res.send(200);
}
}
else {
const response = await dlg.other.request({
@@ -640,10 +811,10 @@ 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 dlg = await this.srf.createUAC(referTo.uri, {localSdp: dlg.local.sdp, headers});
this.uas = dlg;
this.uas.other = this.uac;
+76 -35
View File
@@ -1,7 +1,7 @@
const debug = require('debug')('jambonz:sbc-outbound');
const parseUri = require('drachtio-srf').parseUri;
const Registrar = require('@jambonz/mw-registrar');
const {selectHostPort, makeCallCountKey} = require('./utils');
const {selectHostPort, nudgeCallCounts} = require('./utils');
const FS_UUID_SET_NAME = 'fsUUIDs';
module.exports = (srf, logger, opts) => {
@@ -10,13 +10,15 @@ module.exports = (srf, logger, opts) => {
const registrar = new Registrar(opts);
const {
lookupAccountCapacitiesBySid,
lookupAccountBySid
lookupAccountBySid,
queryCallLimits
} = srf.locals.dbHelpers;
const initLocals = async(req, res, next) => {
req.locals = req.locals || {};
const callId = req.get('Call-ID');
req.locals.account_sid = req.get('X-Account-Sid');
req.locals.application_sid = req.get('X-Application-Sid');
const traceId = req.locals.trace_id = req.get('X-Trace-ID');
req.locals.logger = logger.child({
callId,
@@ -63,7 +65,6 @@ module.exports = (srf, logger, opts) => {
}
}
stats.increment('sbc.invites', ['direction:outbound']);
req.on('cancel', () => {
@@ -76,6 +77,7 @@ module.exports = (srf, logger, opts) => {
try {
req.locals.account = await lookupAccountBySid(req.locals.account_sid);
req.locals.service_provider_sid = req.locals.account.service_provider_sid;
} catch (err) {
req.locals.logger.error({err}, `Error looking up account sid ${req.locals.account_sid}`);
res.send(500);
@@ -85,26 +87,27 @@ module.exports = (srf, logger, opts) => {
};
const checkLimits = async(req, res, next) => {
const {logger, account_sid} = req.locals;
const {writeAlerts, AlertType} = req.srf.locals;
const {logger, account_sid, service_provider_sid, application_sid} = req.locals;
const trackingOn = process.env.JAMBONES_TRACK_ACCOUNT_CALLS ||
process.env.JAMBONES_TRACK_SP_CALLS ||
process.env.JAMBONES_TRACK_APP_CALLS;
if (!process.env.JAMBONES_HOSTING && !trackingOn) {
logger.debug('tracking is off, skipping call limit checks');
return next(); // skip
}
const {writeCallCount, writeCallCountSP, writeCallCountApp, writeAlerts, AlertType} = req.srf.locals;
const key = makeCallCountKey(account_sid);
try {
/* increment the call count */
const calls = await incrKey(key);
debug(`checkLimits: call count is now ${calls}`);
/* decrement count if INVITE is later rejected */
res.once('end', ({status}) => {
res.once('end', async({status}) => {
if (status > 200) {
debug('checkLimits: decrementing call count due to rejection');
decrKey(key)
.then((count) => {
logger.debug({key}, `after rejection there are ${count} active calls for this account`);
debug({key}, `after rejection there are ${count} active calls for this account`);
return;
})
.catch((err) => logger.error({err}, 'checkLimits: decrKey err'));
nudgeCallCounts(logger, {
service_provider_sid,
account_sid,
application_sid
}, decrKey, {writeCallCountSP, writeCallCount, writeCallCountApp})
.catch((err) => logger.error(err, 'Error decrementing call counts'));
const tags = ['accepted:no', `sipStatus:${status}`];
stats.increment('sbc.originations', tags);
}
@@ -114,6 +117,13 @@ module.exports = (srf, logger, opts) => {
}
});
/* increment the call count */
const {callsSP, calls} = await nudgeCallCounts(logger, {
service_provider_sid,
account_sid,
application_sid
}, incrKey, {writeCallCountSP, writeCallCount, writeCallCountApp});
/* compare to account's limit, though avoid db hit when call count is low */
const minLimit = process.env.MIN_CALL_LIMIT ?
parseInt(process.env.MIN_CALL_LIMIT) :
@@ -122,23 +132,54 @@ module.exports = (srf, logger, opts) => {
const capacities = await lookupAccountCapacitiesBySid(account_sid);
const limit = capacities.find((c) => c.category == 'voice_call_session');
if (!limit) {
logger.debug('checkLimits: no call limits specified');
return next();
if (limit) {
const limit_sessions = limit.quantity;
if (calls > limit_sessions) {
logger.info({calls, limit_sessions}, 'checkLimits: limits exceeded');
writeAlerts({
alert_type: AlertType.ACCOUNT_CALL_LIMIT,
service_provider_sid,
account_sid,
count: limit_sessions
}).catch((err) => logger.info({err}, 'checkLimits: error writing alert'));
res.send(503, 'Maximum Calls In Progress');
return req.srf.endSession(req);
}
}
const limit_sessions = limit.quantity;
if (calls > limit_sessions) {
debug(`checkLimits: limits exceeded: call count ${calls}, limit ${limit_sessions}`);
logger.info({calls, limit_sessions}, 'checkLimits: limits exceeded');
writeAlerts({
alert_type: AlertType.CALL_LIMIT,
account_sid,
count: limit_sessions
}).catch((err) => logger.info({err}, 'checkLimits: error writing alert'));
res.send(503, 'Maximum Calls In Progress');
return req.srf.endSession(req);
else if (trackingOn) {
const {account_limit, sp_limit} = await queryCallLimits(service_provider_sid, account_sid);
if (process.env.JAMBONES_TRACK_ACCOUNT_CALLS && account_limit > 0 && calls > account_limit) {
logger.info({calls, account_limit}, 'checkLimits: account limits exceeded');
writeAlerts({
alert_type: AlertType.ACCOUNT_CALL_LIMIT,
service_provider_sid: service_provider_sid,
account_sid,
count: calls
}).catch((err) => logger.info({err}, 'checkLimits: error writing alert'));
res.send(503, 'Max Account Calls In Progress', {
headers: {
'X-Account-Sid': account_sid,
'X-Call-Limit': account_limit
}
});
return req.srf.endSession(req);
}
if (process.env.JAMBONES_TRACK_SP_CALLS && sp_limit > 0 && callsSP > sp_limit) {
logger.info({callsSP, sp_limit}, 'checkLimits: service provider limits exceeded');
writeAlerts({
alert_type: AlertType.SP_CALL_LIMIT,
service_provider_sid: service_provider_sid,
count: callsSP
}).catch((err) => logger.info({err}, 'checkLimits: error writing alert'));
res.send(503, 'Max Service Provider Calls In Progress', {
headers: {
'X-Service-Provider-Sid': service_provider_sid,
'X-Call-Limit': sp_limit
}
});
return req.srf.endSession(req);
}
}
next();
} catch (err) {
+100 -6
View File
@@ -4,9 +4,17 @@ const debug = require('debug')('jambonz:sbc-outbound');
function makeRtpEngineOpts(req, srcIsUsingSrtp, dstIsUsingSrtp, teams = false) {
const from = req.getParsedHeader('from');
const srtpOpts = teams ? srtpCharacteristics['teams'] : srtpCharacteristics['default'];
const dstOpts = dstIsUsingSrtp ? srtpOpts : rtpCharacteristics;
const srcOpts = srcIsUsingSrtp ? srtpOpts : rtpCharacteristics;
const rtpCopy = JSON.parse(JSON.stringify(rtpCharacteristics));
const srtpCopy = JSON.parse(JSON.stringify(srtpCharacteristics));
const srtpOpts = teams ? srtpCopy['teams'] : srtpCopy['default'];
const dstOpts = dstIsUsingSrtp ? srtpOpts : rtpCopy;
const srcOpts = srcIsUsingSrtp ? srtpOpts : rtpCopy;
/* webrtc clients (e.g. sipjs) send DMTF via SIP INFO */
if ((srcIsUsingSrtp || dstIsUsingSrtp) && !teams) {
dstOpts.flags.push('inject DTMF');
srcOpts.flags.push('inject DTMF');
}
const common = {
'call-id': req.get('Call-ID'),
'replace': ['origin', 'session-connection']
@@ -71,7 +79,9 @@ const pingMsTeamsGateways = (logger, srf) => {
});
};
const makeCallCountKey = (sid) => `${sid}:outcalls`;
const makeAccountCallCountKey = (sid) => `outcalls:account:${sid}`;
const makeSPCallCountKey = (sid) => `outcalls:sp:${sid}`;
const makeAppCallCountKey = (sid) => `outcalls:app:${sid}`;
const equalsIgnoreOrder = (a, b) => {
if (a.length !== b.length) return false;
@@ -84,10 +94,94 @@ const equalsIgnoreOrder = (a, b) => {
return true;
};
const systemHealth = async(redisClient, ping, getCount) => {
await Promise.all([redisClient.ping(), ping()]);
return getCount();
};
const createHealthCheckApp = (port, logger) => {
const express = require('express');
const app = express();
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);
});
});
};
const nudgeCallCounts = async(logger, sids, nudgeOperator, writers) => {
const {service_provider_sid, account_sid, application_sid} = sids;
const {writeCallCount, writeCallCountSP, writeCallCountApp} = writers;
const nudges = [];
const writes = [];
logger.debug(sids, 'nudgeCallCounts');
if (process.env.JAMBONES_TRACK_SP_CALLS) {
const key = makeSPCallCountKey(service_provider_sid);
nudges.push(nudgeOperator(key));
}
else {
nudges.push(() => Promise.resolve(null));
}
if (process.env.JAMBONES_TRACK_ACCOUNT_CALLS || process.env.JAMBONES_HOSTING) {
const key = makeAccountCallCountKey(account_sid);
nudges.push(nudgeOperator(key));
}
else {
nudges.push(() => Promise.resolve(null));
}
if (process.env.JAMBONES_TRACK_APP_CALLS && application_sid) {
const key = makeAppCallCountKey(application_sid);
nudges.push(nudgeOperator(key));
}
else {
nudges.push(() => Promise.resolve(null));
}
try {
const [callsSP, calls, callsApp] = await Promise.all(nudges);
logger.debug({
calls, callsSP, callsApp,
service_provider_sid, account_sid, application_sid}, 'call counts after adjustment');
if (process.env.JAMBONES_TRACK_SP_CALLS) {
writes.push(writeCallCountSP({service_provider_sid, calls_in_progress: callsSP}));
}
if (process.env.JAMBONES_TRACK_ACCOUNT_CALLS || process.env.JAMBONES_HOSTING) {
writes.push(writeCallCount({service_provider_sid, account_sid, calls_in_progress: calls}));
}
if (process.env.JAMBONES_TRACK_APP_CALLS && application_sid) {
writes.push(writeCallCountApp({service_provider_sid, account_sid, application_sid, calls_in_progress: callsApp}));
}
/* write the call counts to the database */
Promise.all(writes).catch((err) => logger.error({err}, 'Error writing call counts'));
return {callsSP, calls, callsApp};
} catch (err) {
logger.error(err, 'error incrementing call counts');
}
return {callsSP: null, calls: null, callsApp: null};
};
module.exports = {
makeRtpEngineOpts,
selectHostPort,
pingMsTeamsGateways,
makeCallCountKey,
equalsIgnoreOrder
makeAccountCallCountKey,
makeSPCallCountKey,
equalsIgnoreOrder,
systemHealth,
createHealthCheckApp,
nudgeCallCounts
};
+2220 -2573
View File
File diff suppressed because it is too large Load Diff
+15 -12
View File
@@ -1,6 +1,6 @@
{
"name": "sbc-outbound",
"version": "v0.7.5",
"version": "v0.7.7",
"main": "app.js",
"engines": {
"node": ">= 12.0.0"
@@ -22,29 +22,32 @@
"description": "jambonz session border controller application for outbound calls",
"scripts": {
"start": "node app",
"test": "NODE_ENV=test JAMBONZ_HOSTING=1 JAMBONES_NETWORK_CIDR=127.0.0.1/32 JAMBONES_MYSQL_HOST=127.0.0.1 JAMBONES_MYSQL_USER=jambones_test JAMBONES_MYSQL_PASSWORD=jambones_test JAMBONES_MYSQL_DATABASE=jambones_test JAMBONES_REDIS_HOST=localhost JAMBONES_REDIS_PORT=16379 JAMBONES_TIME_SERIES_HOST=127.0.0.1 JAMBONES_LOGLEVEL=error DRACHTIO_SECRET=cymru DRACHTIO_HOST=127.0.0.1 DRACHTIO_PORT=9060 JAMBONES_RTPENGINES=127.0.0.1:12222 node test/ ",
"test": "NODE_ENV=test HTTP_PORT=3050 JAMBONES_HOSTING=1 JAMBONES_NETWORK_CIDR=127.0.0.1/32 JAMBONES_MYSQL_HOST=127.0.0.1 JAMBONES_MYSQL_USER=jambones_test JAMBONES_MYSQL_PASSWORD=jambones_test JAMBONES_MYSQL_DATABASE=jambones_test JAMBONES_REDIS_HOST=localhost JAMBONES_REDIS_PORT=16379 JAMBONES_TIME_SERIES_HOST=127.0.0.1 JAMBONES_LOGLEVEL=error DRACHTIO_SECRET=cymru DRACHTIO_HOST=127.0.0.1 DRACHTIO_PORT=9060 JAMBONES_RTPENGINES=127.0.0.1:12222 node test/ ",
"coverage": "./node_modules/.bin/nyc --reporter html --report-dir ./coverage npm run test",
"jslint": "eslint app.js lib"
},
"dependencies": {
"@jambonz/db-helpers": "^0.6.17",
"@jambonz/db-helpers": "^0.6.19",
"@jambonz/realtimedb-helpers": "^0.4.35",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/mw-registrar": "0.2.2",
"@jambonz/realtimedb-helpers": "^0.4.24",
"@jambonz/rtpengine-utils": "^0.3.1",
"@jambonz/rtpengine-utils": "^0.3.10",
"@jambonz/siprec-client-utils": "^0.1.4",
"@jambonz/stats-collector": "^0.1.6",
"@jambonz/time-series": "^0.1.9",
"@jambonz/time-series": "^0.2.5",
"cidr-matcher": "^2.1.1",
"debug": "^4.3.3",
"debug": "^4.3.4",
"drachtio-fn-b2b-sugar": "^0.0.12",
"drachtio-srf": "^4.4.59",
"husky": "^7.0.4",
"pino": "^7.4.1"
"drachtio-srf": "^4.5.1",
"express": "^4.18.1",
"pino": "^7.11.0",
"sdp-transform": "^2.14.1"
},
"devDependencies": {
"bent": "^7.3.12",
"eslint": "^7.32.0",
"eslint-plugin-promise": "^5.1.1",
"eslint-plugin-promise": "^5.2.0",
"nyc": "^15.1.0",
"tape": "^5.3.2"
"tape": "^5.5.3"
}
}
+4
View File
@@ -5,3 +5,7 @@ DRACHTIO_SECRET=cymru
JAMBONES_REDIS_HOST=172.39.0.11
JAMBONES_REDIS_PORT=6379
JAMBONES_LOGLEVEL=info
JAMBONES_MYSQL_HOST=172.39.0.2
JAMBONES_MYSQL_USER=jambones_test
JAMBONES_MYSQL_PASSWORD=jambones_test
JAMBONES_MYSQL_DATABASE=jambones_test
+7 -1
View File
@@ -2,7 +2,8 @@ const test = require('tape');
const { output, sippUac } = require('./sipp')('test_sbc-outbound');
const {execSync} = require('child_process');
const debug = require('debug')('jambonz:sbc-outbound');
const consoleLogger = {error: console.error, info: console.log, debug: console.log};
const bent = require('bent');
const getJSON = bent('json');
process.on('unhandledRejection', (reason, p) => {
console.log('Unhandled Rejection at: Promise', p, 'reason:', reason);
@@ -29,6 +30,11 @@ test('sbc-outbound tests', async(t) => {
try {
await connect(srf);
let obj = await getJSON('http://127.0.0.1:3050/');
t.ok(obj.calls === 0, 'HTTP GET / works (current call count)')
obj = await getJSON('http://127.0.0.1:3050/system-health');
t.ok(obj.calls === 0, 'HTTP GET /system-health works (health check)')
/* call to unregistered user */
debug('successfully connected to drachtio server');
await sippUac('uac-pcap-device-404.xml');