Compare commits

..
2 Commits
21 changed files with 1639 additions and 3408 deletions
+1 -1
View File
@@ -8,7 +8,7 @@
"jsx": false,
"modules": false
},
"ecmaVersion": 2020
"ecmaVersion": 2018
},
"plugins": ["promise"],
"rules": {
@@ -11,8 +11,8 @@ jobs:
- uses: actions/checkout@v2
- uses: actions/setup-node@v1
with:
node-version: 14.x
- run: npm ci
node-version: 12
- run: npm install
- run: npm run jslint
- run: npm test
+1 -1
View File
@@ -1,7 +1,7 @@
# Logs
logs
*.log
.vscode/
# Runtime data
pids
*.pid
-4
View File
@@ -1,4 +0,0 @@
#!/bin/sh
. "$(dirname "$0")/_/husky.sh"
npm run jslint
+6 -19
View File
@@ -1,23 +1,10 @@
FROM --platform=linux/amd64 node:16.15.1-alpine as base
RUN apk --update --no-cache add --virtual .builds-deps build-base python3
FROM node:16
WORKDIR /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/
COPY package.json ./
RUN npm install
RUN npm prune
COPY . /opt/app
ARG NODE_ENV
ENV NODE_ENV $NODE_ENV
CMD [ "node", "app.js" ]
CMD [ "npm", "start" ]
+23 -89
View File
@@ -6,28 +6,27 @@ assert.ok(process.env.JAMBONES_MYSQL_HOST &&
assert.ok(process.env.DRACHTIO_PORT || process.env.DRACHTIO_HOST, 'missing DRACHTIO_PORT env var');
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');
assert.ok(process.env.JAMBONES_NETWORK_CIDR, 'missing JAMBONES_NETWORK_CIDR env var');
const Srf = require('drachtio-srf');
const srf = new Srf('sbc-inbound');
const CIDRMatcher = require('cidr-matcher');
const matcher = new CIDRMatcher([process.env.JAMBONES_NETWORK_CIDR]);
const opts = Object.assign({
timestamp: () => {return `, "time": "${new Date().toISOString()}"`;}
}, {level: process.env.JAMBONES_LOGLEVEL || 'info'});
const logger = require('pino')(opts);
const {
writeCallCount,
queryCdrs,
writeCdrs,
writeAlerts,
AlertType
} = require('@jambonz/time-series')(logger, {
host: process.env.JAMBONES_TIME_SERIES_HOST,
port: process.env.JAMBONES_TIME_SERIES_PORT || 8086,
commitSize: 50,
commitInterval: 'test' === process.env.NODE_ENV ? 7 : 20
});
const StatsCollector = require('@jambonz/stats-collector');
const stats = new StatsCollector(logger);
const {equalsIgnoreOrder, createHealthCheckApp, systemHealth} = require('./lib/utils');
const {LifeCycleEvents} = require('./lib/constants');
const setNameRtp = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-rtp`;
const rtpServers = [];
@@ -35,7 +34,6 @@ const setName = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-sip`;
const {
pool,
ping,
lookupAuthHook,
lookupSipGatewayBySignalingAddress,
addSbcAddress,
@@ -51,26 +49,17 @@ const {
database: process.env.JAMBONES_MYSQL_DATABASE,
connectionLimit: process.env.JAMBONES_MYSQL_CONNECTION_LIMIT || 10
}, logger);
const {
client: redisClient,
createSet,
retrieveSet,
addToSet,
removeFromSet,
incrKey,
decrKey} = require('@jambonz/realtimedb-helpers')({
const {createSet, retrieveSet, addToSet, removeFromSet, incrKey, decrKey} = require('@jambonz/realtimedb-helpers')({
host: process.env.JAMBONES_REDIS_HOST || 'localhost',
port: process.env.JAMBONES_REDIS_PORT || 6379
}, logger);
const {getRtpEngine, setRtpEngines} = require('@jambonz/rtpengine-utils')([], logger, {
emitter: stats,
dtmfListenPort: process.env.DTMF_LISTEN_PORT || 22224,
protocol: 'udp'
dtmfListenPort: process.env.DTMF_LISTEN_PORT || 22224
});
srf.locals = {...srf.locals,
stats,
writeCallCount,
queryCdrs,
writeCdrs,
writeAlerts,
@@ -79,7 +68,6 @@ srf.locals = {...srf.locals,
getRtpEngine,
dbHelpers: {
pool,
ping,
lookupAuthHook,
lookupSipGatewayBySignalingAddress,
lookupAccountByPhoneNumber,
@@ -117,13 +105,7 @@ const {
} = require('./lib/middleware')(srf, logger);
const CallSession = require('./lib/call-session');
if (process.env.DRACHTIO_HOST && !process.env.K8S) {
const CIDRMatcher = require('cidr-matcher');
const cidrs = process.env.JAMBONES_NETWORK_CIDR
.split(',')
.map((s) => s.trim());
const matcher = new CIDRMatcher(cidrs);
if (process.env.DRACHTIO_HOST) {
srf.connect({host: process.env.DRACHTIO_HOST, port: process.env.DRACHTIO_PORT, secret: process.env.DRACHTIO_SECRET });
srf.on('connect', (err, hp) => {
if (err) return this.logger.error({err}, 'Error connecting to drachtio server');
@@ -149,9 +131,6 @@ if (process.env.DRACHTIO_HOST && !process.env.K8S) {
});
}
else {
srf.on('listening', () => {
logger.info(`listening in outbound mode on port ${process.env.DRACHTIO_PORT}`);
});
srf.listen({port: process.env.DRACHTIO_PORT, secret: process.env.DRACHTIO_SECRET});
}
if (process.env.NODE_ENV === 'test') {
@@ -184,78 +163,33 @@ srf.use((req, res, next, err) => {
res.send(500);
});
if (process.env.K8S || process.env.HTTP_PORT) {
const PORT = process.env.HTTP_PORT || 3000;
const healthCheck = require('@jambonz/http-health-check');
/* update call stats periodically */
setInterval(() => {
stats.gauge('sbc.sip.calls.count', activeCallIds.size, ['direction:inbound']);
}, 20000);
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(() => {
stats.gauge('sbc.sip.calls.count', activeCallIds.size, ['direction:inbound']);
}, 20000);
}
const lookupRtpServiceEndpoints = (lookup, serviceName) => {
logger.debug(`dns lookup for ${serviceName}..`);
lookup(serviceName, {family: 4, all: true}, (err, addresses) => {
if (err) {
logger.error({err}, `Error looking up ${serviceName}`);
return;
}
logger.debug({addresses, rtpServers}, `dns lookup for ${serviceName} returned`);
const addrs = addresses.map((a) => a.address);
if (!equalsIgnoreOrder(addrs, rtpServers)) {
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}`));
}
});
const arrayCompare = (a, b) => {
if (a.length !== b.length) return false;
const uniqueValues = new Set([...a, ...b]);
for (const v of uniqueValues) {
const aCount = a.filter((e) => e === v).length;
const bCount = b.filter((e) => e === v).length;
if (aCount !== bCount) return false;
}
return true;
};
if (process.env.K8S_RTPENGINE_SERVICE_NAME) {
/* poll dns for endpoints every so often */
const arr = /^(.*):(\d+)$/.exec(process.env.K8S_RTPENGINE_SERVICE_NAME);
const svc = arr[1];
logger.info(`rtpengine(s) will be found at dns name: ${svc}`);
const {lookup} = require('dns');
lookupRtpServiceEndpoints(lookup, svc);
setInterval(lookupRtpServiceEndpoints.bind(null, lookup, svc), process.env.RTPENGINE_DNS_POLL_INTERVAL || 10000);
}
else if (process.env.JAMBONES_RTPENGINES) {
/* static list of rtpengines */
/* update rtpengines periodically */
if (process.env.JAMBONES_RTPENGINES) {
setRtpEngines([process.env.JAMBONES_RTPENGINES]);
}
else {
/* poll redis periodically for rtpengines that have registered via OPTIONS ping */
const getActiveRtpServers = async() => {
try {
const set = await retrieveSet(setNameRtp);
const newArray = Array.from(set);
logger.debug({newArray, rtpServers}, 'getActiveRtpServers');
if (!equalsIgnoreOrder(newArray, rtpServers)) {
if (!arrayCompare(newArray, rtpServers)) {
logger.info({newArray}, 'resetting active rtpengines');
setRtpEngines(newArray.map((a) => `${a}:${process.env.RTPENGINE_PORT || 22222}`));
rtpServers.length = 0;
@@ -265,11 +199,11 @@ else {
logger.error({err}, 'Error setting new rtpengines');
}
};
setInterval(() => {
getActiveRtpServers();
}, 30000);
getActiveRtpServers();
}
const {lifecycleEmitter} = require('./lib/autoscale-manager')(logger);
-29
View File
@@ -1,29 +0,0 @@
#!/usr/bin/env node
const bent = require('bent');
const getJSON = bent('json');
const PORT = process.env.HTTP_PORT || 3000;
const sleep = (ms) => {
return new Promise((resolve) => setTimeout(resolve, ms));
};
(async function() {
try {
do {
const obj = await getJSON(`http://127.0.0.1:${PORT}/`);
const {calls} = obj;
if (calls === 0) {
console.log('no calls on the system, we can exit');
process.exit(0);
}
else {
console.log(`waiting for ${calls} to exit..`);
}
await sleep(10000);
} while (1);
} catch (err) {
console.error(err, 'Error querying health endpoint');
process.exit(-1);
}
})();
+15 -152
View File
@@ -1,5 +1,4 @@
const Emitter = require('events');
const SrsClient = require('./srs-client');
const {makeRtpEngineOpts, SdpWantsSrtp, makeCallCountKey} = require('./utils');
const {forwardInDialogRequests} = require('drachtio-fn-b2b-sugar');
const {parseUri, stringifyUri, SipError} = require('drachtio-srf');
@@ -63,10 +62,7 @@ class CallSession extends Emitter {
blockDTMF,
unblockDTMF,
subscribeDTMF,
unsubscribeDTMF,
subscribeRequest,
subscribeAnswer,
unsubscribe
unsubscribeDTMF
} = engine;
this.offer = offer;
this.answer = answer;
@@ -77,9 +73,6 @@ class CallSession extends Emitter {
this.unblockDTMF = unblockDTMF;
this.subscribeDTMF = subscribeDTMF;
this.unsubscribeDTMF = unsubscribeDTMF;
this.subscribeRequest = subscribeRequest;
this.subscribeAnswer = subscribeAnswer;
this.unsubscribe = unsubscribe;
const featureServer = await this.getFeatureServer();
if (!featureServer) {
@@ -178,7 +171,7 @@ class CallSession extends Emitter {
'-Session-Expires',
'-X-Subspace-Forwarded-For'
],
proxyResponseHeaders: ['all', '-X-Trace-ID'],
proxyResponseHeaders: ['all'],
localSdpB: response.sdp,
localSdpA: async(sdp, res) => {
this.rtpEngineOpts.uac.tag = res.getParsedHeader('To').params.tag;
@@ -204,19 +197,18 @@ class CallSession extends Emitter {
this._setHandlers({uas, uac});
return;
} catch (err) {
this.rtpEngineResource.destroy().catch((err) => this.logger.info({err}, 'Error destroying rtpe after failure'));
this.rtpEngineResource.destroy();
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');
return this.emit('failed');
}
else if (err.message !== 'call canceled') {
this.logger.error(err, 'unexpected error routing inbound call');
}
this.srf.endSession(this.req);
}
}
@@ -231,13 +223,6 @@ class CallSession extends Emitter {
this.rtpEngineResource.destroy().catch((err) => {});
this.activeCallIds.delete(callId);
if (dlg.other && dlg.other.connected) dlg.other.destroy().catch((e) => {});
if (this.srsClient) {
this.srsClient.stop();
this.srsClient = null;
}
this.srf.endSession(this.req);
});
//re-invite
@@ -254,30 +239,22 @@ class CallSession extends Emitter {
this.req.locals.cdr = {
...this.req.locals.cdr,
answered: true,
answered_at: callStart,
trace_id: uac.res?.get('X-Trace-ID') || '00000000000000000000000000000000'
answered_at: callStart
};
}
this.uas = uas;
this.uac = uac;
[uas, uac].forEach((dlg) => {
dlg.on('destroy', async() => {
const other = dlg.other;
dlg.on('destroy', () => {
this.logger.info('call ended with normal termination');
this.rtpEngineResource.destroy().catch((err) => {});
this.activeCallIds.delete(this.req.get('Call-ID'));
try {
await other.destroy();
} catch (err) {}
dlg.other.destroy().catch((e) => {});
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
if (process.env.JAMBONES_HOSTING || process.env.JAMBONES_TRACK_ACCOUNT_CALLS) {
const {account_sid} = this.req.locals;
if (process.env.JAMBONES_HOSTING) {
this.decrKey(this.callCountKey)
.then((count) => {
this.logger.info(
{key: this.callCountKey},
`after hangup there are ${count} active calls for this account`);
return this.req.srf.locals.writeCallCount({account_sid, calls_in_progress: count});
})
.then((count) => this.logger.debug({key: this.callCountKey},
`after hangup there are ${count} active calls for this account`))
.catch((err) => this.logger.error({err}, 'Error decrementing call count'));
}
@@ -295,19 +272,6 @@ class CallSession extends Emitter {
trunk
}).catch((err) => this.logger.error({err}, 'Error writing cdr for completed call'));
}
/* de-link the 2 Dialogs for GC */
dlg.removeAllListeners();
other.removeAllListeners();
dlg.other = null;
other.other = null;
if (this.srsClient) {
this.srsClient.stop();
this.srsClient = null;
}
this.logger.info(`call ended with normal termination, there are ${this.activeCallIds.size} active`);
this.srf.endSession(this.req);
});
});
@@ -431,6 +395,7 @@ Duration=${payload.duration} `
direction,
sdp: req.body,
};
//if (reason && opts.flags && !opts.flags.includes('reset')) opts.flags.push('reset');
let response = await this.offer(opts);
if ('ok' !== response.result) {
@@ -467,7 +432,6 @@ Duration=${payload.duration} `
async _onInfo(dlg, req, res) {
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;
try {
if (dlg.type === 'uac' && req.has('X-Reason')) {
const reason = req.get('X-Reason');
@@ -477,100 +441,16 @@ Duration=${payload.duration} `
'from-tag': fromTag
};
this.logger.info(`_onInfo: got request ${reason}`);
res.send(200);
if (reason.startsWith('mute')) {
const response = Promise.all([this.blockMedia(opts), this.blockDTMF(opts)]);
res.send(200);
this.logger.info({response}, `_onInfo: response to rtpengine command for ${reason}`);
}
else if (reason.startsWith('unmute')) {
const response = Promise.all([this.unblockMedia(opts), this.unblockDTMF(opts)]);
res.send(200);
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 a 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,
originalInvite: this.req,
callingNumber: this.req.callingNumber,
calledNumber: this.req.calledNumber,
srsUrl,
srsRecordingId,
callSid,
accountSid,
applicationSid,
rtpEngineOpts: this.rtpEngineOpts,
fromTag,
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 {
const immutableHdrs = ['via', 'from', 'to', 'call-id', 'cseq', 'max-forwards', 'content-length'];
@@ -586,10 +466,6 @@ Duration=${payload.duration} `
res.send(response.status, {headers: responseHeaders, body: response.body});
}
} catch (err) {
if (this.srsClient) {
this.srsClient = null;
}
res.send(500);
this.logger.info({err}, `Error handing INFO request on ${dlg.type} leg`);
}
}
@@ -653,7 +529,6 @@ Duration=${payload.duration} `
this.rtpEngineResource.destroy();
this.activeCallIds.delete(this.req.get('Call-ID'));
uac.other.destroy();
this.srf.endSession(this.req);
});
// now we can destroy the old dialog
dlg.destroy().catch(() => {});
@@ -739,20 +614,8 @@ Duration=${payload.duration} `
}
}
else {
/* REFER coming in from a sip device, forward to feature server */
try {
const response = await dlg.other.request({
method: 'REFER',
headers: {
'Refer-To': req.get('Refer-To'),
'Referred-By': req.get('Referred-By'),
'User-Agent': req.get('User-Agent')
}
});
res.send(response.status, response.reason);
} catch (err) {
this.logger.error({err}, 'CallSession:_onRefer: error handling incoming REFER');
}
// TODO: forward on to feature server
res.send(501);
}
}
+2 -15
View File
@@ -44,10 +44,6 @@ SELECT * FROM phone_numbers
WHERE number = ?
AND voip_carrier_sid = ?`;
const sqlQueryAllDidsForCarrier = `
SELECT * FROM phone_numbers
WHERE voip_carrier_sid = ?`;
const sqlSelectOutboundGatewayForCarrier = `
SELECT ipv4, port, e164_leading_plus
FROM sip_gateways sg, voip_carriers vc
@@ -85,18 +81,9 @@ module.exports = (srf, logger) => {
const did = normalizeDID(req.calledNumber);
try {
/* straight DID match */
const [r] = await pp.query(sqlQueryApplicationByDid, [did, voip_carrier_sid]);
if (r.length) return r[0].application_sid;
/* wildcard / regex match */
const [r2] = await pp.query(sqlQueryAllDidsForCarrier, [voip_carrier_sid]);
const match = r2
.filter((o) => o.number.match(/\D/)) // look at anything with non-digit characters
.sort((a, b) => b.number.length - a.number.length) // prefer longest match
.find((o) => did.match(new RegExp(o.number.endsWith('*') ? `${o.number.slice(0, -1)}\\d*` : o.number)));
if (match) return match.application_sid;
return null;
if (0 === r.length) return null;
return r[0].application_sid;
} catch (err) {
logger.error({err}, 'getApplicationForDidAndCarrier');
}
+6 -15
View File
@@ -1,8 +1,4 @@
const setName = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-fs`;
const assert = require('assert');
assert.ok(!process.env.K8S || process.env.K8S_FEATURE_SERVER_SERVICE_NAME,
'when running in Kubernetes, an env var K8S_FEATURE_SERVER_SERVICE_NAME is required');
module.exports = (srf, logger) => {
const {retrieveSet, createSet} = srf.locals.realtimeDbHelpers;
@@ -14,18 +10,13 @@ module.exports = (srf, logger) => {
return async() => {
try {
if (process.env.K8S) {
return process.env.K8S_FEATURE_SERVER_SERVICE_NAME;
}
else {
const fs = await retrieveSet(setName);
if (0 === fs.length) {
logger.info('No available feature servers to handle incoming call');
return;
}
logger.debug({fs}, `retrieved ${setName}`);
return fs[idx++ % fs.length];
const fs = await retrieveSet(setName);
if (0 === fs.length) {
logger.info('No available feature servers to handle incoming call');
return;
}
logger.debug({fs}, `retrieved ${setName}`);
return fs[idx++ % fs.length];
} catch (err) {
logger.error({err}, `Error retrieving ${setName}`);
}
+8 -19
View File
@@ -152,8 +152,7 @@ module.exports = function(srf, logger) {
const app = await lookupAppByTeamsTenant(uri.host);
if (!app) {
stats.increment('sbc.terminations', ['sipStatus:404']);
res.send(404, {headers: {'X-Reason': 'no configured application'}});
return req.srf.endSession(req);
return res.send(404, {headers: {'X-Reason': 'no configured application'}});
}
req.locals = {
@@ -172,8 +171,7 @@ module.exports = function(srf, logger) {
const account = await lookupAccountBySipRealm(uri.host);
if (!account) {
stats.increment('sbc.terminations', ['sipStatus:404']);
res.send(404);
return req.srf.endSession(req);
return res.send(404);
}
/* if this is a dedicated SBC (static IP) only take calls for that account's sip realm */
@@ -182,8 +180,7 @@ module.exports = function(srf, logger) {
`identifyAccount: static IP for ${process.env.SBC_ACCOUNT_SID} but call for ${account.account_sid}`);
stats.increment('sbc.terminations', ['sipStatus:404']);
delete req.locals.cdr;
res.send(404);
return req.srf.endSession(req);
return res.send(404);
}
req.locals = {
account_sid: account.account_sid,
@@ -220,11 +217,11 @@ module.exports = function(srf, logger) {
};
const checkLimits = async(req, res, next) => {
if (!process.env.JAMBONES_HOSTING && !process.env.JAMBONES_TRACK_ACCOUNT_CALLS) return next(); // skip
if (!process.env.JAMBONES_HOSTING) return next(); // skip
const {incrKey, decrKey} = req.srf.locals.realtimeDbHelpers;
const {logger, account_sid} = req.locals;
const {writeCallCount, writeAlerts, AlertType} = req.srf.locals;
const {writeAlerts, AlertType} = req.srf.locals;
assert(account_sid);
const key = makeCallCountKey(account_sid);
@@ -233,11 +230,10 @@ module.exports = function(srf, logger) {
if (status > 200) {
decrKey(key)
.then((count) => {
logger.info({key}, `after rejection there are ${count} active calls for this account`);
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 count;
return;
})
.then((count) => writeCallCount({account_sid, calls_in_progress: count}))
.catch((err) => logger.error({err}, 'checkLimits: decrKey err'));
}
});
@@ -245,10 +241,6 @@ module.exports = function(srf, logger) {
try {
/* increment the call count */
const calls = await incrKey(key);
writeCallCount({account_sid, calls_in_progress: calls})
.then(() => logger.info(`checkLimits: after incrementing there are ${calls} active calls for this account`))
.catch((err) => logger.error({err}, 'checkLimits: error writing call count'));
if (!process.env.JAMBONES_HOSTING) return next();
/* compare to account's limit, though avoid db hit when call count is low */
const minLimit = process.env.MIN_CALL_LIMIT ?
@@ -269,15 +261,13 @@ module.exports = function(srf, logger) {
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);
return res.send(503, 'Maximum Calls In Progress');
}
next();
} catch (err) {
stats.increment('sbc.terminations', ['sipStatus:500']);
logger.error({err}, 'error checking limits error for inbound call');
res.send(500);
req.srf.endSession(req);
}
};
@@ -290,7 +280,6 @@ module.exports = function(srf, logger) {
stats.increment('sbc.terminations', ['sipStatus:500']);
logger.error(err, `${req.get('Call-ID')} Error looking up related info for inbound call`);
res.send(500);
req.srf.endSession(req);
}
};
-249
View File
@@ -1,249 +0,0 @@
const Emitter = require('events');
const assert = require('assert');
const transform = require('sdp-transform');
const { v4: uuidv4 } = require('uuid');
const createMultipartSdp = (sdp, {
originalInvite,
srsRecordingId,
callSid,
accountSid,
applicationSid,
sipCallId,
aorFrom,
aorTo,
callingNumber,
calledNumber
}) => {
const sessionId = uuidv4();
const uuidStream1 = uuidv4();
const uuidStream2 = uuidv4();
const participant1 = uuidv4();
const participant2 = uuidv4();
const sipSessionId = originalInvite.get('Call-ID');
const {originator = 'unknown', carrier = 'unknown'} = originalInvite.locals;
const x = `--uniqueBoundary
Content-Disposition: session;handling=required
Content-Type: application/sdp
--sdp-placeholder--
--uniqueBoundary
Content-Disposition: recording-session
Content-Type: application/rs-metadata+xml
<?xml version="1.0" encoding="UTF-8"?>
<recording xmlns="urn:ietf:params:xml:ns:recording:1">
<datamode>complete</datamode>
<session session_id="${sessionId}">
<sipSessionID>${sipSessionId}</sipSessionID>
</session>
<extensiondata xmlns:jb="http://jambonz.org/siprec">
<jb:callsid>${callSid}</jb:callsid>
<jb:accountsid>${accountSid}</jb:accountsid>
<jb:applicationsid>${applicationSid}</jb:applicationsid>
<jb:recordingid>${srsRecordingId}</jb:recordingid>
<jb:originationsource>${originator}</jb:originationsource>
<jb:carrier>${carrier}</jb:carrier>
</extensiondata>
<participant participant_id="${participant1}">
<nameID aor="${aorFrom}">
<name>${callingNumber}</name>
</nameID>
</participant>
<participantsessionassoc participant_id="${participant1}" session_id="${sessionId}">
</participantsessionassoc>
<stream stream_id="${uuidStream1}" session_id="${sessionId}">
<label>1</label>
</stream>
<participant participant_id="${participant2}">
<nameID aor="${aorTo}">
<name>${calledNumber}</name>
</nameID>
</participant>
<participantsessionassoc participant_id="${participant2}" session_id="${sessionId}">
</participantsessionassoc>
<stream stream_id="${uuidStream2}" session_id="${sessionId}">
<label>2</label>
</stream>
<participantstreamassoc participant_id="${participant1}">
<send>${uuidStream1}</send>
<recv>${uuidStream2}</recv>
</participantstreamassoc>
<participantstreamassoc participant_id="${participant2}">
<send>${uuidStream2}</send>
<recv>${uuidStream1}</recv>
</participantstreamassoc>
</recording>`
.replace(/\n/g, '\r\n')
.replace('--sdp-placeholder--', sdp);
return `${x}\r\n`;
};
class SrsClient extends Emitter {
constructor(logger, opts) {
super();
const {
srf,
originalInvite,
calledNumber,
callingNumber,
srsUrl,
srsRecordingId,
callSid,
accountSid,
applicationSid,
srsDestUserName,
rtpEngineOpts,
//fromTag,
toTag,
aorFrom,
aorTo,
subscribeRequest,
subscribeAnswer,
del,
blockMedia,
unblockMedia,
unsubscribe
} = opts;
this.logger = logger;
this.srf = srf;
this.originalInvite = originalInvite;
this.callingNumber = callingNumber;
this.calledNumber = calledNumber;
this.subscribeRequest = subscribeRequest;
this.subscribeAnswer = subscribeAnswer;
this.del = del;
this.blockMedia = blockMedia;
this.unblockMedia = unblockMedia;
this.unsubscribe = unsubscribe;
this.srsUrl = srsUrl;
this.srsRecordingId = srsRecordingId;
this.callSid = callSid;
this.accountSid = accountSid;
this.applicationSid = applicationSid;
this.srsDestUserName = srsDestUserName;
this.rtpEngineOpts = rtpEngineOpts;
this.sipRecFromTag = toTag;
this.aorFrom = aorFrom;
this.aorTo = aorTo;
/* state */
this.activated = false;
this.paused = false;
}
async start() {
assert(!this.activated);
const opts = {
'call-id': this.rtpEngineOpts.common['call-id'],
'from-tag': this.sipRecFromTag
};
let response = await this.subscribeRequest({...opts, label: '1', flags: ['all'], interface: 'public'});
if (response.result !== 'ok') {
this.logger.error({response}, 'SrsClient:start error calling subscribe request');
throw new Error('error calling subscribe request');
}
this.siprecToTag = response['to-tag'];
const parsed = transform.parse(response.sdp);
parsed.name = 'jambonz SRS';
parsed.media[0].label = '1';
parsed.media[1].label = '2';
this.sdpOffer = transform.write(parsed);
const sdp = createMultipartSdp(this.sdpOffer, {
originalInvite: this.originalInvite,
srsRecordingId: this.srsRecordingId,
callSid: this.callSid,
accountSid: this.accountSid,
applicationSid: this.applicationSid,
calledNumber: this.calledNumber,
callingNumber: this.callingNumber,
aorFrom: this.aorFrom,
aorTo: this.aorTo
});
this.logger.info({response}, `SrsClient: sending SDP ${sdp}`);
/* */
try {
this.uac = await this.srf.createUAC(this.srsUrl, {
headers: {
'Content-Type': 'multipart/mixed;boundary=uniqueBoundary',
},
localSdp: sdp
});
} catch (err) {
this.logger.info({err}, `Error sending SIPREC INVITE to ${this.srsUrl}`);
throw err;
}
this.logger.info({sdp: this.uac.remote.sdp}, `SrsClient:start - successfully connected to SRS ${this.srsUrl}`);
response = await this.subscribeAnswer({
...opts,
sdp: this.uac.remote.sdp,
'to-tag': response['to-tag'],
label: '2'
});
if (response.result !== 'ok') {
this.logger.error({response}, 'SrsClient:start error calling subscribe answer');
throw new Error('error calling subscribe answer');
}
this.activated = true;
this.logger.info('successfully established siprec connection');
return true;
}
async stop() {
assert(this.activated);
const opts = {
'call-id': this.rtpEngineOpts.common['call-id'],
'from-tag': this.sipRecFromTag
};
this.del(opts)
//.then((response) => this.logger.debug({response}, 'Successfully stopped siprec media'))
.catch((err) => this.logger.info({err}, 'Error deleting siprec media session'));
this.uac.destroy().catch(() => {});
this.activated = false;
return true;
}
async pause() {
assert(!this.paused);
const opts = {
'call-id': this.rtpEngineOpts.common['call-id'],
'from-tag': this.sipRecFromTag
};
try {
await this.blockMedia(opts);
await this.uac.modify(this.sdpOffer.replace(/sendonly/g, 'inactive'));
this.paused = true;
return true;
} catch (err) {
this.logger.info({err}, 'Error pausing siprec media session');
}
return false;
}
async resume() {
assert(this.paused);
const opts = {
'call-id': this.rtpEngineOpts.common['call-id'],
'from-tag': this.sipRecFromTag
};
try {
await this.blockMedia(opts);
await this.uac.modify(this.sdpOffer);
} catch (err) {
this.logger.info({err}, 'Error resuming siprec media session');
}
return true;
}
}
module.exports = SrsClient;
+1 -35
View File
@@ -45,45 +45,11 @@ const normalizeDID = (tel) => {
return arr ? arr[1] : tel;
};
const equalsIgnoreOrder = (a, b) => {
if (a.length !== b.length) return false;
const uniqueValues = new Set([...a, ...b]);
for (const v of uniqueValues) {
const aCount = a.filter((e) => e === v).length;
const bCount = b.filter((e) => e === v).length;
if (aCount !== bCount) return false;
}
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);
});
});
};
module.exports = {
isWSS,
SdpWantsSrtp,
getAppserver,
makeRtpEngineOpts,
makeCallCountKey,
normalizeDID,
equalsIgnoreOrder,
systemHealth,
createHealthCheckApp
normalizeDID
};
-24
View File
@@ -1,24 +0,0 @@
#!/bin/sh
TCP_SERVER_PORT="${DRACHTIO_PORT:-4000}"
nc -v -z localhost $TCP_SERVER_PORT
# if last command exited with non zero
if [ $? != 0 ]
then
exit 1
fi
HTTP_SERVER_PORT="${HTTP_PORT:-3000}"
printf 'GET /system-health HTTP/1.1\r\nHost: localhost\r\n\r\n' | nc -v localhost 3000 | grep calls
# grep will automatically exit with 1 if string is not matched, however, will leave that call there in case
# we pivot to pipe to dev/null
if [ $? != 0 ]
then
exit 1
fi
exit 0
+1546 -2574
View File
File diff suppressed because it is too large Load Diff
+21 -22
View File
@@ -1,9 +1,9 @@
{
"name": "sbc-inbound",
"version": "v0.7.5",
"version": "v0.7.1",
"main": "app.js",
"engines": {
"node": ">= 12.0.0"
"node": ">= 10.16.0"
},
"keywords": [
"sip",
@@ -20,35 +20,34 @@
},
"scripts": {
"start": "node app",
"test": "NODE_ENV=test HTTP_PORT=3050 JAMBONES_NETWORK_CIDR='127.0.0.1/32' JAMBONES_HOSTING=1 SBC_ACCOUNT_SID=ed649e33-e771-403a-8c99-1780eabbc803 JAMBONES_TIME_SERIES_HOST=127.0.0.1 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_LOGLEVEL=info DRACHTIO_SECRET=cymru DRACHTIO_HOST=127.0.0.1 DRACHTIO_PORT=9060 JAMBONES_RTPENGINES=127.0.0.1:12222 JAMBONES_FEATURE_SERVERS=172.38.0.11 node test/ ",
"test": "NODE_ENV=test JAMBONES_NETWORK_CIDR='127.0.0.1/32' JAMBONES_HOSTING=1 SBC_ACCOUNT_SID=ed649e33-e771-403a-8c99-1780eabbc803 JAMBONES_TIME_SERIES_HOST=127.0.0.1 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_LOGLEVEL=error DRACHTIO_SECRET=cymru DRACHTIO_HOST=127.0.0.1 DRACHTIO_PORT=9060 JAMBONES_RTPENGINES=127.0.0.1:12222 JAMBONES_FEATURE_SERVERS=172.38.0.11 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.18",
"@jambonz/http-authenticator": "^0.2.1",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/realtimedb-helpers": "^0.4.29",
"@jambonz/rtpengine-utils": "^0.3.1",
"@jambonz/stats-collector": "^0.1.6",
"@jambonz/time-series": "^0.1.9",
"aws-sdk": "^2.1152.0",
"@jambonz/db-helpers": "^0.6.12",
"@jambonz/http-authenticator": "^0.2.0",
"@jambonz/realtimedb-helpers": "^0.4.8",
"@jambonz/rtpengine-utils": "^0.1.17",
"@jambonz/stats-collector": "^0.1.5",
"@jambonz/time-series": "^0.1.5",
"aws-sdk": "^2.848.0",
"bent": "^7.3.12",
"express": "^4.17.1",
"verify-aws-sns-signature": "^0.0.6",
"xml2js": "^0.4.23",
"cidr-matcher": "^2.1.1",
"debug": "^4.3.4",
"debug": "^4.3.1",
"drachtio-fn-b2b-sugar": "0.0.12",
"drachtio-srf": "^4.5.0",
"express": "^4.18.1",
"pino": "^7.11.0",
"sdp-transform": "^2.14.1",
"uuid": "^8.3.2",
"verify-aws-sns-signature": "^0.0.7",
"xml2js": "^0.4.23"
"drachtio-srf": "^4.4.59",
"pino": "^6.8.0",
"rtpengine-client": "^0.2.0"
},
"devDependencies": {
"eslint": "^7.32.0",
"eslint-plugin-promise": "^4.3.1",
"clear-module": "^4.1.1",
"eslint": "^7.15.0",
"eslint-plugin-promise": "^4.2.1",
"nyc": "^15.1.0",
"tape": "^4.15.1"
"tape": "^4.13.3"
}
}
-7
View File
@@ -56,12 +56,5 @@ values ('888a5339-c62c-4075-9e19-f4de70a96597', '999c1452-620d-4195-9f19-c9814ef
insert into phone_numbers (phone_number_sid, number, voip_carrier_sid, account_sid)
values ('999a5339-c62c-4075-9e19-f4de70a96597', '16173333456', '287c1452-620d-4195-9f19-c9814ef90d78', 'ed649e33-e771-403a-8c99-1780eabbc803');
insert into phone_numbers (phone_number_sid, number, voip_carrier_sid, account_sid)
values ('29543d4e-d959-4a25-836a-cde7161cd7d5', '1508222*', '287c1452-620d-4195-9f19-c9814ef90d78', 'ed649e33-e771-403a-8c99-1780eabbc803');
insert into phone_numbers (phone_number_sid, number, voip_carrier_sid, account_sid)
values ('dddd5c34-feae-4d70-98af-bb4d1f8dc965', '1508*', '287c1452-620d-4195-9f19-c9814ef90d78', 'ed649e33-e771-403a-8c99-1780eabbc803');
insert into phone_numbers (phone_number_sid, number, voip_carrier_sid, account_sid)
values ('d458bf7a-bcea-47b2-ac96-66dfc9c5c220', '150822233*', '287c1452-620d-4195-9f19-c9814ef90d78', 'ed649e33-e771-403a-8c99-1780eabbc803');
insert into phone_numbers (phone_number_sid, number, voip_carrier_sid, account_sid)
values ('f7ad205d-b92f-4363-8160-f8b5216b40d3', '15083871234', '287c1452-620d-4195-9f19-c9814ef90d78', 'd7cc37cb-d152-49ef-a51b-485f6e917089');
+1 -2
View File
@@ -10,7 +10,6 @@ networks:
services:
mysql:
image: mysql:5.7
platform: linux/x86_64
ports:
- "3306:3306"
environment:
@@ -72,7 +71,7 @@ services:
ipv4_address: 172.38.0.14
influxdb:
image: influxdb:1.8
image: influxdb:1.8-alpine
ports:
- "8086:8086"
networks:
-118
View File
@@ -1,118 +0,0 @@
<?xml version="1.0" encoding="ISO-8859-1" ?>
<!DOCTYPE scenario SYSTEM "sipp.dtd">
<!-- This program is free software; you can redistribute it and/or -->
<!-- modify it under the terms of the GNU General Public License as -->
<!-- published by the Free Software Foundation; either version 2 of the -->
<!-- License, or (at your option) any later version. -->
<!-- -->
<!-- This program is distributed in the hope that it will be useful, -->
<!-- but WITHOUT ANY WARRANTY; without even the implied warranty of -->
<!-- MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the -->
<!-- GNU General Public License for more details. -->
<!-- -->
<!-- You should have received a copy of the GNU General Public License -->
<!-- along with this program; if not, write to the -->
<!-- Free Software Foundation, Inc., -->
<!-- 59 Temple Place, Suite 330, Boston, MA 02111-1307 USA -->
<!-- -->
<!-- Sipp 'uac' scenario with pcap (rtp) play -->
<!-- -->
<scenario name="UAC with media">
<!-- In client mode (sipp placing calls), the Call-ID MUST be -->
<!-- generated by sipp. To do so, use [call_id] keyword. -->
<send retrans="500">
<![CDATA[
INVITE sip:+15082223333@jambonz.org SIP/2.0
Via: SIP/2.0/[transport] [local_ip]:[local_port];branch=[branch]
From: sipp <sip:sipp@[local_ip]:[local_port]>;tag=[pid]SIPpTag09[call_number]
To: <sip:15082223333@jambonz.org>
Call-ID: [call_id]
CSeq: 1 INVITE
Contact: sip:sipp@[local_ip]:[local_port]
Max-Forwards: 70
Subject: uac-pcap-carrier-success
Content-Type: application/sdp
Content-Length: [len]
v=0
o=user1 53655765 2353687637 IN IP[local_ip_type] [local_ip]
s=-
c=IN IP[local_ip_type] [local_ip]
t=0 0
m=audio [auto_media_port] RTP/AVP 8 101
a=rtpmap:8 PCMA/8000
a=rtpmap:101 telephone-event/8000
a=fmtp:101 0-11,16
]]>
</send>
<recv response="100" optional="true">
</recv>
<recv response="180" optional="true">
</recv>
<!-- By adding rrs="true" (Record Route Sets), the route sets -->
<!-- are saved and used for following messages sent. Useful to test -->
<!-- against stateful SIP proxies/B2BUAs. -->
<recv response="200" rtd="true" crlf="true">
</recv>
<!-- Packet lost can be simulated in any send/recv message by -->
<!-- by adding the 'lost = "10"'. Value can be [1-100] percent. -->
<send>
<![CDATA[
ACK sip:15082223333@jambonz.org SIP/2.0
Via: SIP/2.0/[transport] [local_ip]:[local_port];branch=[branch]
From: sipp <sip:sipp@[local_ip]:[local_port]>;tag=[pid]SIPpTag09[call_number]
To: <sip:15082223333@jambonz.org>[peer_tag_param]
Call-ID: [call_id]
CSeq: 1 ACK
Max-Forwards: 70
Subject: uac-pcap-carrier-success
Content-Length: 0
]]>
</send>
<!-- Play a pre-recorded PCAP file (RTP stream) -->
<nop>
<action>
<exec play_pcap_audio="pcap/g711a.pcap"/>
</action>
</nop>
<!-- Pause briefly -->
<pause milliseconds="3000"/>
<!-- The 'crlf' option inserts a blank line in the statistics report. -->
<send retrans="500">
<![CDATA[
BYE sip:15082223333@jambonz.org SIP/2.0
Via: SIP/2.0/[transport] [local_ip]:[local_port];branch=[branch]
From: sipp <sip:sipp@[local_ip]:[local_port]>;tag=[pid]SIPpTag09[call_number]
To: <sip:15082223333@jambonz.org>[peer_tag_param]
Call-ID: [call_id]
CSeq: 2 BYE
Subject: uac-pcap-carrier-success
Content-Length: 0
]]>
</send>
<recv response="200" crlf="true">
</recv>
<!-- definition of the response time repartition table (unit is ms) -->
<ResponseTimeRepartition value="10, 20, 30, 40, 50, 100, 150, 200"/>
<!-- definition of the call length repartition table (unit is ms) -->
<CallLengthRepartition value="10, 50, 100, 500, 1000, 5000, 10000"/>
</scenario>
+1 -18
View File
@@ -124,24 +124,7 @@
]]>
</send>
<recv request="BYE">
</recv>
<send next="2">
<![CDATA[
SIP/2.0 200 OK
[last_Via:]
[last_From:]
[last_To:]
[last_Call-ID:]
[last_CSeq:]
Contact: <sip:[local_ip]:[local_port];transport=[transport]>
Content-Length: 0
]]>
</send>
<label id="2"/>
</scenario>
+5 -13
View File
@@ -1,7 +1,8 @@
const test = require('tape');
const { sippUac } = require('./sipp')('test_sbc-inbound');
const bent = require('bent');
const getJSON = bent('json');
const { output, sippUac } = require('./sipp')('test_sbc-inbound');
const debug = require('debug')('drachtio:sbc-inbound');
const clearModule = require('clear-module');
const consoleLogger = {error: console.error, info: console.log, debug: console.log};
process.on('unhandledRejection', (reason, p) => {
console.log('Unhandled Rejection at: Promise', p, 'reason:', reason);
@@ -27,21 +28,12 @@ test('incoming call 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)')
await sippUac('uac-pcap-carrier-success.xml', '172.38.0.20');
t.pass('incoming call from carrier completed successfully');
await sippUac('uac-pcap-pbx-success.xml', '172.38.0.21');
t.pass('incoming call from account-level carrier completed successfully');
await sippUac('uac-did-regex-match.xml', '172.38.0.20');
t.pass('incoming call matched by trailing wildcard *');
await sippUac('uac-pcap-device-success.xml', '172.38.0.30');
t.pass('incoming call from authenticated device completed successfully');
@@ -63,7 +55,7 @@ test('incoming call tests', async(t) => {
await waitFor(10);
const res = await queryCdrs({account_sid: 'ed649e33-e771-403a-8c99-1780eabbc803'});
console.log(`cdrs: ${JSON.stringify(res)}`);
t.ok(7 === res.total, 'successfully wrote 7 cdrs for calls');
t.ok(6 === res.total, 'successfully wrote 6 cdrs for calls');
srf.disconnect();
t.end();