Compare commits

..
16 changed files with 1113 additions and 3639 deletions
-51
View File
@@ -1,51 +0,0 @@
name: Docker
on:
push:
# Publish `master` as Docker `latest` image.
branches:
- main
# Publish `v1.2.3` tags as releases.
tags:
- v*
env:
IMAGE_NAME: sbc-inbound
jobs:
push:
runs-on: ubuntu-latest
if: github.event_name == 'push'
steps:
- uses: actions/checkout@v2
- name: Build image
run: docker build . --file Dockerfile --tag $IMAGE_NAME
- name: Log into registry
run: echo "${{ secrets.GITHUB_TOKEN }}" | docker login ghcr.io -u ${{ github.actor }} --password-stdin
- name: Push image
run: |
IMAGE_ID=ghcr.io/${{ github.repository_owner }}/$IMAGE_NAME
# Change all uppercase to lowercase
IMAGE_ID=$(echo $IMAGE_ID | tr '[A-Z]' '[a-z]')
# Strip git ref prefix from version
VERSION=$(echo "${{ github.ref }}" | sed -e 's,.*/\(.*\),\1,')
# Strip "v" prefix from tag name
[[ "${{ github.ref }}" == "refs/tags/"* ]] && VERSION=$(echo $VERSION | sed -e 's/^v//')
# Use Docker `latest` tag convention
[ "$VERSION" == "main" ] && VERSION=latest
echo IMAGE_ID=$IMAGE_ID
echo VERSION=$VERSION
docker tag $IMAGE_NAME $IMAGE_ID:$VERSION
docker push $IMAGE_ID:$VERSION
@@ -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
+8 -2
View File
@@ -1,10 +1,16 @@
FROM node:17.4-slim
FROM node:alpine as builder
RUN apk update && apk add --no-cache python make g++
WORKDIR /opt/app/
COPY package.json ./
RUN npm install
RUN npm prune
FROM node:alpine as app
WORKDIR /opt/app
COPY . /opt/app
COPY --from=builder /opt/app/node_modules ./node_modules
ARG NODE_ENV
ENV NODE_ENV $NODE_ENV
CMD [ "npm", "start" ]
CMD [ "npm", "start" ]
+25 -80
View File
@@ -6,9 +6,11 @@ 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'});
@@ -20,13 +22,11 @@ const {
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} = require('./lib/utils');
const {LifeCycleEvents} = require('./lib/constants');
const setNameRtp = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-rtp`;
const rtpServers = [];
@@ -54,11 +54,7 @@ const {createSet, retrieveSet, addToSet, removeFromSet, incrKey, decrKey} = requ
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'
});
const {getRtpEngine, setRtpEngines} = require('@jambonz/rtpengine-utils')([], logger, {emitter: stats});
srf.locals = {...srf.locals,
stats,
queryCdrs,
@@ -84,18 +80,7 @@ srf.locals = {...srf.locals,
retrieveSet
}
};
const {
wasOriginatedFromCarrier,
getApplicationForDidAndCarrier,
getOutboundGatewayForRefer
} = require('./lib/db-utils')(srf, logger);
srf.locals = {
...srf.locals,
wasOriginatedFromCarrier,
getApplicationForDidAndCarrier,
getOutboundGatewayForRefer,
getFeatureServer: require('./lib/fs-tracking')(srf, logger)
};
srf.locals.getFeatureServer = require('./lib/fs-tracking')(srf, logger);
const activeCallIds = srf.locals.activeCallIds;
const {
@@ -106,17 +91,12 @@ 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');
logger.info(`connected to drachtio listening on ${hp}`);
if (process.env.SBC_ACCOUNT_SID) return;
const hostports = hp.split(',');
for (const hp of hostports) {
@@ -124,12 +104,11 @@ if (process.env.DRACHTIO_HOST && !process.env.K8S) {
if (arr && 'udp' === arr[1] && !matcher.contains(arr[2])) {
logger.info(`adding sbc public address to database: ${arr[2]}`);
srf.locals.sipAddress = arr[2];
if (!process.env.SBC_ACCOUNT_SID) addSbcAddress(arr[2]);
addSbcAddress(arr[2]);
}
else if (arr && 'tcp' === arr[1] && matcher.contains(arr[2])) {
const hostport = `${arr[2]}:${arr[3]}`;
logger.info(`adding sbc private address to redis: ${hostport}`);
srf.locals.privateSipAddress = hostport;
srf.locals.addToRedis = () => addToSet(setName, hostport);
srf.locals.removeFromRedis = () => removeFromSet(setName, hostport);
srf.locals.addToRedis();
@@ -138,7 +117,6 @@ if (process.env.DRACHTIO_HOST && !process.env.K8S) {
});
}
else {
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') {
@@ -171,58 +149,33 @@ srf.use((req, res, next, err) => {
res.send(500);
});
if (process.env.K8S) {
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});
}
if ('test' !== process.env.NODE_ENV) {
/* update call stats periodically */
setInterval(() => {
stats.gauge('sbc.sip.calls.count', activeCallIds.size, ['direction:inbound']);
}, 20000);
}
/* 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;
@@ -232,11 +185,11 @@ else {
logger.error({err}, 'Error setting new rtpengines');
}
};
setInterval(() => {
getActiveRtpServers();
}, 30000);
getActiveRtpServers();
}
const {lifecycleEmitter} = require('./lib/autoscale-manager')(logger);
@@ -251,12 +204,4 @@ setInterval(async() => {
}
}, 20000);
process.on('SIGUSR2', handle.bind(null, removeFromSet, setName));
process.on('SIGTERM', handle.bind(null, removeFromSet, setName));
function handle(removeFromSet, setName, signal) {
logger.info(`got signal ${signal}, removing ${srf.locals.privateSipAddress} from set ${setName}`);
removeFromSet(setName, srf.locals.privateSipAddress);
}
module.exports = {srf, 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);
}
})();
+1 -1
View File
@@ -3,6 +3,6 @@
"DTLS": "off",
"SDES": "off",
"ICE": "remove",
"flags": ["media handover", "port latching"],
"flags": ["media handover"],
"rtcp-mux": ["accept"]
}
+2 -2
View File
@@ -3,14 +3,14 @@
"transport-protocol": "UDP/TLS/RTP/SAVPF",
"ICE": "force",
"SDES": "off",
"flags": ["generate mid", "SDES-no", "media handover", "port latching"],
"flags": ["generate mid", "SDES-no", "media handover"],
"rtcp-mux": ["require"]
},
"teams": {
"transport-protocol": "RTP/SAVP",
"ICE": "force",
"SDES": "off",
"flags": ["generate mid", "SDES-no", "media handover", "port latching"],
"flags": ["generate mid", "SDES-no", "media handover"],
"rtcp-mux": ["accept"]
}
}
+18 -172
View File
@@ -1,7 +1,7 @@
const Emitter = require('events');
const {makeRtpEngineOpts, SdpWantsSrtp, makeCallCountKey} = require('./utils');
const {forwardInDialogRequests} = require('drachtio-fn-b2b-sugar');
const {parseUri, stringifyUri, SipError} = require('drachtio-srf');
const {parseUri, 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';
@@ -39,10 +39,6 @@ class CallSession extends Emitter {
return !!this.req.locals.msTeamsTenantFqdn;
}
get privateSipAddress() {
return this.srf.locals.privateSipAddress;
}
async connect() {
this.logger.info('inbound call accepted for routing');
const engine = this.getRtpEngine();
@@ -53,26 +49,10 @@ class CallSession extends Emitter {
return this.res.send(480);
}
debug(`got engine: ${JSON.stringify(engine)}`);
const {
offer,
answer,
del,
blockMedia,
unblockMedia,
blockDTMF,
unblockDTMF,
subscribeDTMF,
unsubscribeDTMF
} = engine;
const {offer, answer, del} = engine;
this.offer = offer;
this.answer = answer;
this.del = del;
this.blockMedia = blockMedia;
this.unblockMedia = unblockMedia;
this.blockDTMF = blockDTMF;
this.unblockDTMF = unblockDTMF;
this.subscribeDTMF = subscribeDTMF;
this.unsubscribeDTMF = unsubscribeDTMF;
const featureServer = await this.getFeatureServer();
if (!featureServer) {
@@ -118,22 +98,15 @@ class CallSession extends Emitter {
}
// now send the INVITE in towards the feature servers
let headers = {
const headers = {
'From': createBLegFromHeader(this.req),
'To': this.req.get('To'),
'X-Account-Sid': this.req.locals.account_sid,
'X-CID': this.req.get('Call-ID'),
'X-Forwarded-For': `${this.req.source_address}`
'X-Forwarded-For': `${this.req.source_address}:${this.req.source_port}`
};
if (this.privateSipAddress) headers = {...headers, Contact: `<sip:${this.privateSipAddress}>`};
const responseHeaders = {};
if (this.req.locals.carrier) {
Object.assign(headers, {
'X-Originating-Carrier': this.req.locals.carrier,
'X-Voip-Carrier-Sid': this.req.locals.voip_carrier_sid
});
}
if (this.req.locals.carrier) Object.assign(headers, {'X-Originating-Carrier': this.req.locals.carrier});
if (this.req.locals.msTeamsTenantFqdn) {
Object.assign(headers, {'X-MS-Teams-Tenant-FQDN': this.req.locals.msTeamsTenantFqdn});
@@ -169,7 +142,7 @@ class CallSession extends Emitter {
'-Max-Forwards',
'-Record-Route',
'-Session-Expires',
'-X-Subspace-Forwarded-For'
'-X-Forwarded-For'
],
proxyResponseHeaders: ['all'],
localSdpB: response.sdp,
@@ -197,7 +170,7 @@ 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) {
@@ -209,22 +182,17 @@ class CallSession extends Emitter {
else if (err.message !== 'call canceled') {
this.logger.error(err, 'unexpected error routing inbound call');
}
this.srf.endSession(this.req);
}
}
_setDlgHandlers(dlg) {
const {callId} = dlg.sip;
this.activeCallIds.set(callId, this);
this.subscribeDTMF(this.logger, callId, this.rtpEngineOpts.uas.tag,
this._onDTMF.bind(this));
this.activeCallIds.set(this.req.get('Call-ID'), this);
dlg.on('destroy', () => {
debug('call ended with normal termination');
this.logger.info('call ended with normal termination');
this.rtpEngineResource.destroy().catch((err) => {});
this.activeCallIds.delete(callId);
this.activeCallIds.delete(this.req.get('Call-ID'));
if (dlg.other && dlg.other.connected) dlg.other.destroy().catch((e) => {});
this.srf.endSession(this.req);
});
//re-invite
@@ -252,7 +220,7 @@ class CallSession extends Emitter {
this.rtpEngineResource.destroy().catch((err) => {});
this.activeCallIds.delete(this.req.get('Call-ID'));
dlg.other.destroy().catch((e) => {});
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
if (process.env.JAMBONES_HOSTING) {
this.decrKey(this.callCountKey)
.then((count) => this.logger.debug({key: this.callCountKey},
@@ -274,52 +242,16 @@ class CallSession extends Emitter {
trunk
}).catch((err) => this.logger.error({err}, 'Error writing cdr for completed call'));
}
this.srf.endSession(this.req);
});
});
this.subscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag,
this._onDTMF.bind(this, uac));
uas.on('modify', this._onReinvite.bind(this, uas));
uac.on('modify', this._onReinvite.bind(this, uac));
uac.on('refer', this._onFeatureServerTransfer.bind(this, uac));
uas.on('refer', this._onRefer.bind(this, uas));
uas.on('info', this._onInfo.bind(this, uas));
uac.on('info', this._onInfo.bind(this, uac));
// default forwarding of other request types
forwardInDialogRequests(uas, ['notify', 'options', 'message']);
}
async _onDTMF(dlg, payload) {
this.logger.info({payload}, '_onDTMF');
try {
let dtmf;
switch (payload.event) {
case 10:
dtmf = '*';
break;
case 11:
dtmf = '#';
break;
default:
dtmf = '' + payload.event;
break;
}
await dlg.request({
method: 'INFO',
headers: {
'Content-Type': 'application/dtmf-relay'
},
body: `Signal=${dtmf}
Duration=${payload.duration} `
});
} catch (err) {
this.logger.info({err}, 'Error sending INFO application/dtmf-relay');
}
forwardInDialogRequests(uas, ['info', 'notify', 'options', 'message']);
}
/**
@@ -329,7 +261,7 @@ Duration=${payload.duration} `
*/
async replaces(req, res) {
try {
let opts = Object.assign({}, this.rtpEngineOpts.uas.mediaOpts, {sdp: req.body});
let opts = Object.assign(this.rtpEngineOpts.offer, {sdp: req.body});
let response = await this.offer(opts);
if ('ok' !== response.result) {
res.send(488);
@@ -337,8 +269,8 @@ Duration=${payload.duration} `
}
this.logger.info({opts, response}, 'sent offer for reinvite to rtpengine');
const sdp = await this.uac.modify(response.sdp);
opts = Object.assign({}, this.rtpEngineOpts.uac.mediaOpts, {sdp, 'to-tag': this.toTag});
Object.assign(this.rtpEngineOpts.uas.mediaOpts, {'to-tag': this.toTag});
opts = Object.assign(this.rtpEngineOpts.answer, {sdp, 'to-tag': this.toTag});
Object.assign(this.rtpEngineOpts.offer, {'to-tag': this.toTag});
response = await this.answer(opts);
if ('ok' !== response.result) {
res.send(488);
@@ -356,8 +288,6 @@ Duration=${payload.duration} `
});
}
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
const uas = await this.srf.createUAS(req, res, {
localSdp: response.sdp,
headers
@@ -384,7 +314,6 @@ Duration=${payload.duration} `
res.send(200, {body: dlg.local.sdp});
return;
}
const reason = req.get('X-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;
@@ -398,23 +327,13 @@ 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) {
res.send(488);
throw new Error(`_onReinvite: rtpengine failed: offer: ${JSON.stringify(response)}`);
}
/* if this is a re-invite from the FS to change media anchoring, avoid sending the reinvite out */
let sdp;
if (reason && dlg.type === 'uac' && ['release-media', 'anchor-media'].includes(reason)) {
this.logger.info({response}, `got a reinvite from FS to ${reason}`);
sdp = dlg.other.remote.sdp;
}
else {
sdp = await dlg.other.modify(response.sdp);
}
const sdp = await dlg.other.modify(response.sdp);
opts = {
...this.rtpEngineOpts.common,
...answerMedia,
@@ -433,86 +352,15 @@ Duration=${payload.duration} `
}
}
async _onInfo(dlg, req, res) {
const fromTag = dlg.type === 'uas' ? this.rtpEngineOpts.uas.tag : this.rtpEngineOpts.uac.tag;
try {
if (dlg.type === 'uac' && req.has('X-Reason')) {
const reason = req.get('X-Reason');
const opts = {
...this.rtpEngineOpts.common,
flags: ['reset'],
'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)]);
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)]);
this.logger.info({response}, `_onInfo: response to rtpengine command for ${reason}`);
}
}
else {
const immutableHdrs = ['via', 'from', 'to', 'call-id', 'cseq', 'max-forwards', 'content-length'];
const headers = {};
Object.keys(req.headers).forEach((h) => {
if (!immutableHdrs.includes(h)) headers[h] = req.headers[h];
});
const response = await dlg.other.request({method: 'INFO', headers, body: req.body});
const responseHeaders = {};
if (response.has('Content-Type')) {
Object.assign(responseHeaders, {'Content-Type': response.get('Content-Type')});
}
res.send(response.status, {headers: responseHeaders, body: response.body});
}
} catch (err) {
this.logger.info({err}, `Error handing INFO request on ${dlg.type} leg`);
}
}
async _onFeatureServerTransfer(dlg, req, res) {
try {
const referTo = req.getParsedHeader('Refer-To');
const uri = parseUri(referTo.uri);
this.logger.info({uri, referTo, headers: req.headers}, 'received REFER from feature server');
this.logger.info({uri, referTo}, 'received REFER from feature server');
const arr = /context-(.*)/.exec(uri.user);
if (!arr) {
/* call transfer requested */
const {gateway} = this.req.locals;
const referredBy = req.getParsedHeader('Referred-By');
if (!referredBy) return res.send(400);
const u = parseUri(referredBy.uri);
let selectedGateway = false;
let e164 = false;
if (gateway) {
/* host of Refer-to to an outbound gateway */
const gw = await this.srf.locals.getOutboundGatewayForRefer(gateway.voip_carrier_sid);
if (gw) {
selectedGateway = true;
e164 = gw.e164_leading_plus;
uri.host = gw.ipv4;
uri.port = gw.port;
}
}
if (!selectedGateway) {
uri.host = this.req.source_address;
uri.port = this.req.source_port;
}
if (e164 && !uri.user.startsWith('+')) {
uri.user = `+${uri.user}`;
}
const response = await this.uas.request({
method: 'REFER',
headers: {
'Refer-To': stringifyUri(uri),
'Referred-By': stringifyUri(u)
}
});
return res.send(response.status);
this.logger.info(`invalid Refer-To header: ${referTo.uri}`);
return res.send(501);
}
res.send(202);
@@ -532,7 +380,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(() => {});
@@ -604,7 +451,6 @@ Duration=${payload.duration} `
// successfully connected
this.logger.info('successfully connected new call leg for REFER');
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
this.referInvite = null;
sendNotify(this.uas, '200 OK');
this.uas.destroy();
+5 -56
View File
@@ -11,15 +11,6 @@ WHERE acc.sip_realm = ?
AND vc.account_sid = acc.account_sid
AND sg.voip_carrier_sid = vc.voip_carrier_sid`;
const sqlSelectAllCarriersForSPByRealm =
`SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.account_sid,
vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask
FROM sip_gateways sg, voip_carriers vc, accounts acc
WHERE acc.sip_realm = ?
AND vc.service_provider_sid = acc.service_provider_sid
AND vc.account_sid IS NULL
AND sg.voip_carrier_sid = vc.voip_carrier_sid`;
const sqlSelectAllGatewaysForSP =
`SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.service_provider_sid,
vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask
@@ -44,13 +35,6 @@ SELECT * FROM phone_numbers
WHERE number = ?
AND voip_carrier_sid = ?`;
const sqlSelectOutboundGatewayForCarrier = `
SELECT ipv4, port, e164_leading_plus
FROM sip_gateways sg, voip_carriers vc
WHERE sg.voip_carrier_sid = ?
AND sg.voip_carrier_sid = vc.voip_carrier_sid
AND outbound = 1`;
const gatewayMatchesSourceAddress = (source_address, gw) => {
if (32 === gw.netmask && gw.ipv4 === source_address) return true;
if (gw.netmask < 32) {
@@ -64,19 +48,6 @@ module.exports = (srf, logger) => {
const {pool} = srf.locals.dbHelpers;
const pp = pool.promise();
const getOutboundGatewayForRefer = async(voip_carrier_sid) => {
try {
const [r] = await pp.query(sqlSelectOutboundGatewayForCarrier, [voip_carrier_sid]);
if (0 === r.length) return null;
/* if multiple, prefer a DNS name */
const hasDns = r.find((row) => row.ipv4.match(/^[A-Za-z]/));
return hasDns || r[0];
} catch (err) {
logger.error({err}, 'getOutboundGatewayForRefer');
}
};
const getApplicationForDidAndCarrier = async(req, voip_carrier_sid) => {
const did = normalizeDID(req.calledNumber);
@@ -102,11 +73,7 @@ module.exports = (srf, logger) => {
const [r] = await pp.query(sqlCarriersForAccountBySid,
[process.env.SBC_ACCOUNT_SID, req.source_address, req.source_port]);
if (0 === r.length) return failure;
return {
fromCarrier: true,
gateway: r[0],
account_sid: process.env.SBC_ACCOUNT_SID
};
return {fromCarrier: true, gateway: r[0], account_sid: process.env.SBC_ACCOUNT_SID};
}
else {
/* we may have a carrier at the service provider level */
@@ -121,23 +88,8 @@ module.exports = (srf, logger) => {
const [r] = await pp.query(sql);
if (0 === r.length) {
/* came from a carrier, but number is not provisioned..
check if we only have a single account, otherwise we have no
way of knowing which account this is for
*/
const [r] = await pp.query('SELECT count(*) as count from accounts where service_provider_sid = ?',
matches[0].service_provider_sid);
if (r[0].count === 0) return {fromCarrier: true};
else {
const [accounts] = await pp.query('SELECT * from accounts where service_provider_sid = ?',
matches[0].service_provider_sid);
return {
fromCarrier: true,
gateway: matches[0],
account_sid: accounts[0].account_sid,
account: accounts[0]
};
}
/* came from a carrier, but number is not provisioned */
return {fromCarrier: true};
}
const gateway = matches.find((m) => m.voip_carrier_sid === r[0].voip_carrier_sid);
const [accounts] = await pp.query(sqlAccountBySid, r[0].account_sid);
@@ -155,9 +107,7 @@ module.exports = (srf, logger) => {
}
/* get all the carriers and gateways for the account owning this sip realm */
const [gwAcc] = await pp.query(sqlSelectAllCarriersForAccountByRealm, uri.host);
const [gwSP] = gwAcc.length ? [[]] : await pp.query(sqlSelectAllCarriersForSPByRealm, uri.host);
const gw = gwAcc.concat(gwSP);
const [gw] = await pp.query(sqlSelectAllCarriersForAccountByRealm, uri.host);
const selected = gw.find(gatewayMatchesSourceAddress.bind(null, req.source_address));
if (selected) {
const [a] = await pp.query(sqlAccountByRealm, uri.host);
@@ -175,7 +125,6 @@ module.exports = (srf, logger) => {
return {
wasOriginatedFromCarrier,
getApplicationForDidAndCarrier,
getOutboundGatewayForRefer
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}`);
}
+23 -28
View File
@@ -68,19 +68,28 @@ module.exports = function(srf, logger) {
blacklistUnknownRealms: true,
emitter: new AuthOutcomeReporter(stats)
});
const {wasOriginatedFromCarrier, getApplicationForDidAndCarrier} = require('./db-utils')(srf, logger);
const initLocals = (req, res, next) => {
req.locals = req.locals || {};
/* check if forwarded by a proxy that applied an X-Forwarded-For Header */
if (req.has('X-Forwarded-For') || req.has('X-Subspace-Forwarded-For')) {
const original_source_address = req.get('X-Forwarded-For') || req.get('X-Subspace-Forwarded-For');
logger.info({
callId: req.get('Call-ID'),
original_source_address,
proxy_source_address: req.source_address,
}, 'overwriting source address for proxied SIP INVITE');
req.source_address = original_source_address;
if (req.has('X-Forwarded-For')) {
const arr = /^([^:]*)(?::(.*))?$/.exec(req.get('X-Forwarded-For'));
if (arr) {
const original_source_address = arr[1];
const original_source_port = arr[2] || req.source_port;
logger({
callId: req.get('Call-ID'),
original_source_address,
original_source_port,
proxy_source_address: req.source_address,
proxy_source_port: req.source_port
}, 'overwriting source address and port for proxied SIP INVITE');
req.source_address = original_source_address;
req.source_port = original_source_port;
}
}
req.locals.cdr = initCdr(req);
const callId = req.get('Call-ID');
@@ -111,14 +120,8 @@ module.exports = function(srf, logger) {
const identifyAccount = async(req, res, next) => {
try {
const {wasOriginatedFromCarrier, getApplicationForDidAndCarrier} = req.srf.locals;
const {
fromCarrier,
gateway,
account_sid,
application_sid,
account
} = await wasOriginatedFromCarrier(req);
const {fromCarrier, gateway, account_sid, application_sid, account} = await wasOriginatedFromCarrier(req);
/**
* calls come from 3 sources:
* (1) A carrier
@@ -137,8 +140,6 @@ module.exports = function(srf, logger) {
req.locals = {
originator: 'trunk',
carrier: gateway.name,
gateway,
voip_carrier_sid: gateway.voip_carrier_sid,
application_sid: sid || gateway.application_sid,
account_sid,
account,
@@ -152,8 +153,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 +172,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 +181,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,
@@ -264,15 +262,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);
}
};
@@ -285,7 +281,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);
}
};
+1 -13
View File
@@ -45,23 +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;
};
module.exports = {
isWSS,
SdpWantsSrtp,
getAppserver,
makeRtpEngineOpts,
makeCallCountKey,
normalizeDID,
equalsIgnoreOrder
normalizeDID
};
+1002 -3166
View File
File diff suppressed because it is too large Load Diff
+18 -19
View File
@@ -1,9 +1,9 @@
{
"name": "sbc-inbound",
"version": "v0.7.3",
"version": "0.3.6",
"main": "app.js",
"engines": {
"node": ">= 12.0.0"
"node": ">= 10.16.0"
},
"keywords": [
"sip",
@@ -20,34 +20,33 @@
},
"scripts": {
"start": "node app",
"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/ ",
"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=debug 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.16",
"@jambonz/db-helpers": "^0.6.12",
"@jambonz/http-authenticator": "^0.2.0",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/realtimedb-helpers": "^0.4.24",
"@jambonz/rtpengine-utils": "^0.3.1",
"@jambonz/stats-collector": "^0.1.6",
"@jambonz/time-series": "^0.1.6",
"aws-sdk": "^2.1036.0",
"@jambonz/realtimedb-helpers": "^0.4.8",
"@jambonz/rtpengine-utils": "^0.1.12",
"@jambonz/stats-collector": "^0.1.5",
"@jambonz/time-series": "^0.1.5",
"aws-sdk": "^2.848.0",
"bent": "^7.3.12",
"cidr-matcher": "^2.1.1",
"debug": "^4.3.3",
"drachtio-fn-b2b-sugar": "0.0.12",
"drachtio-srf": "^4.4.59",
"express": "^4.17.1",
"pino": "^7.4.1",
"rtpengine-client": "^0.2.0",
"verify-aws-sns-signature": "^0.0.6",
"xml2js": "^0.4.23"
"xml2js": "^0.4.23",
"cidr-matcher": "^2.1.1",
"debug": "^4.3.1",
"drachtio-fn-b2b-sugar": "0.0.12",
"drachtio-srf": "^4.4.49",
"pino": "^6.8.0",
"rtpengine-client": "^0.2.0"
},
"devDependencies": {
"clear-module": "^4.1.1",
"eslint": "^7.32.0",
"eslint-plugin-promise": "^4.3.1",
"eslint": "^7.15.0",
"eslint-plugin-promise": "^4.2.1",
"nyc": "^15.1.0",
"tape": "^4.13.3"
}
+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: