mirror of
https://github.com/jambonz/sbc-inbound.git
synced 2026-10-04 02:04:22 +00:00
Compare commits
7
Commits
v0.7.1
..
v0.7.2-rc2
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
fcf90916e1 | ||
|
|
4a0aa79456 | ||
|
|
1cdaae7773 | ||
|
|
82ac86389d | ||
|
|
e9184ad208 | ||
|
|
804bb890b5 | ||
|
|
a892a87eb5 |
+1
-1
@@ -1,4 +1,4 @@
|
||||
FROM node:16
|
||||
FROM node:17-slim
|
||||
WORKDIR /opt/app/
|
||||
COPY package.json ./
|
||||
RUN npm install
|
||||
|
||||
@@ -6,11 +6,9 @@ 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, 'missing JAMBONES_NETWORK_CIDR env var');
|
||||
assert.ok(process.env.JAMBONES_NETWORK_CIDR || process.env.K8S, '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'});
|
||||
@@ -22,6 +20,7 @@ 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
|
||||
});
|
||||
@@ -56,7 +55,8 @@ const {createSet, retrieveSet, addToSet, removeFromSet, incrKey, decrKey} = requ
|
||||
|
||||
const {getRtpEngine, setRtpEngines} = require('@jambonz/rtpengine-utils')([], logger, {
|
||||
emitter: stats,
|
||||
dtmfListenPort: process.env.DTMF_LISTEN_PORT || 22224
|
||||
dtmfListenPort: process.env.DTMF_LISTEN_PORT || 22224,
|
||||
protocol: process.env.RTPENGINE_NG_PROTOCOL || (process.env.K8S ? 'ws' : 'udp')
|
||||
});
|
||||
srf.locals = {...srf.locals,
|
||||
stats,
|
||||
@@ -105,7 +105,13 @@ const {
|
||||
} = require('./lib/middleware')(srf, logger);
|
||||
const CallSession = require('./lib/call-session');
|
||||
|
||||
if (process.env.DRACHTIO_HOST) {
|
||||
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);
|
||||
|
||||
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');
|
||||
@@ -131,6 +137,7 @@ if (process.env.DRACHTIO_HOST) {
|
||||
});
|
||||
}
|
||||
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') {
|
||||
@@ -163,6 +170,11 @@ srf.use((req, res, next, err) => {
|
||||
res.send(500);
|
||||
});
|
||||
|
||||
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});
|
||||
|
||||
/* update call stats periodically */
|
||||
setInterval(() => {
|
||||
stats.gauge('sbc.sip.calls.count', activeCallIds.size, ['direction:inbound']);
|
||||
@@ -179,11 +191,13 @@ const arrayCompare = (a, b) => {
|
||||
return true;
|
||||
};
|
||||
|
||||
/* update rtpengines periodically */
|
||||
if (process.env.JAMBONES_RTPENGINES) {
|
||||
setRtpEngines([process.env.JAMBONES_RTPENGINES]);
|
||||
const serviceName = process.env.JAMBONES_RTPENGINES || process.env.K8S_RTPENGINE_SERVICE_NAME;
|
||||
if (serviceName) {
|
||||
logger.info(`rtpengine(s) will be found at: ${serviceName}`);
|
||||
setRtpEngines([serviceName]);
|
||||
}
|
||||
else {
|
||||
/* update rtpengines periodically */
|
||||
const getActiveRtpServers = async() => {
|
||||
try {
|
||||
const set = await retrieveSet(setNameRtp);
|
||||
@@ -199,11 +213,11 @@ else {
|
||||
logger.error({err}, 'Error setting new rtpengines');
|
||||
}
|
||||
};
|
||||
|
||||
setInterval(() => {
|
||||
getActiveRtpServers();
|
||||
}, 30000);
|
||||
getActiveRtpServers();
|
||||
|
||||
}
|
||||
|
||||
const {lifecycleEmitter} = require('./lib/autoscale-manager')(logger);
|
||||
|
||||
Executable
+29
@@ -0,0 +1,29 @@
|
||||
#!/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);
|
||||
}
|
||||
})();
|
||||
+9
-5
@@ -209,6 +209,7 @@ class CallSession extends Emitter {
|
||||
else if (err.message !== 'call canceled') {
|
||||
this.logger.error(err, 'unexpected error routing inbound call');
|
||||
}
|
||||
this.srf.endSession(this.req);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -223,6 +224,7 @@ class CallSession extends Emitter {
|
||||
this.rtpEngineResource.destroy().catch((err) => {});
|
||||
this.activeCallIds.delete(callId);
|
||||
if (dlg.other && dlg.other.connected) dlg.other.destroy().catch((e) => {});
|
||||
this.srf.endSession(this.req);
|
||||
});
|
||||
|
||||
//re-invite
|
||||
@@ -272,6 +274,7 @@ class CallSession extends Emitter {
|
||||
trunk
|
||||
}).catch((err) => this.logger.error({err}, 'Error writing cdr for completed call'));
|
||||
}
|
||||
this.srf.endSession(this.req);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -326,7 +329,7 @@ Duration=${payload.duration} `
|
||||
*/
|
||||
async replaces(req, res) {
|
||||
try {
|
||||
let opts = Object.assign(this.rtpEngineOpts.offer, {sdp: req.body});
|
||||
let opts = Object.assign({}, this.rtpEngineOpts.uas.mediaOpts, {sdp: req.body});
|
||||
let response = await this.offer(opts);
|
||||
if ('ok' !== response.result) {
|
||||
res.send(488);
|
||||
@@ -334,8 +337,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.answer, {sdp, 'to-tag': this.toTag});
|
||||
Object.assign(this.rtpEngineOpts.offer, {'to-tag': this.toTag});
|
||||
opts = Object.assign({}, this.rtpEngineOpts.uac.mediaOpts, {sdp, 'to-tag': this.toTag});
|
||||
Object.assign(this.rtpEngineOpts.uas.mediaOpts, {'to-tag': this.toTag});
|
||||
response = await this.answer(opts);
|
||||
if ('ok' !== response.result) {
|
||||
res.send(488);
|
||||
@@ -395,7 +398,7 @@ Duration=${payload.duration} `
|
||||
direction,
|
||||
sdp: req.body,
|
||||
};
|
||||
if (reason) opts.flags.push('reset');
|
||||
//if (reason && opts.flags && !opts.flags.includes('reset')) opts.flags.push('reset');
|
||||
|
||||
let response = await this.offer(opts);
|
||||
if ('ok' !== response.result) {
|
||||
@@ -406,7 +409,7 @@ Duration=${payload.duration} `
|
||||
/* if this is a re-invite from the FS to change media anchoring, avoid sending the reinvite out */
|
||||
let sdp;
|
||||
if (reason && dlg.type === 'uac' && ['release-media', 'anchor-media'].includes(reason)) {
|
||||
this.logger.info(`got a reinvite from FS to ${reason}`);
|
||||
this.logger.info({response}, `got a reinvite from FS to ${reason}`);
|
||||
sdp = dlg.other.remote.sdp;
|
||||
}
|
||||
else {
|
||||
@@ -529,6 +532,7 @@ 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(() => {});
|
||||
|
||||
+15
-6
@@ -1,4 +1,8 @@
|
||||
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;
|
||||
@@ -10,13 +14,18 @@ module.exports = (srf, logger) => {
|
||||
|
||||
return async() => {
|
||||
try {
|
||||
const fs = await retrieveSet(setName);
|
||||
if (0 === fs.length) {
|
||||
logger.info('No available feature servers to handle incoming call');
|
||||
return;
|
||||
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];
|
||||
}
|
||||
logger.debug({fs}, `retrieved ${setName}`);
|
||||
return fs[idx++ % fs.length];
|
||||
} catch (err) {
|
||||
logger.error({err}, `Error retrieving ${setName}`);
|
||||
}
|
||||
|
||||
+10
-4
@@ -152,7 +152,8 @@ module.exports = function(srf, logger) {
|
||||
const app = await lookupAppByTeamsTenant(uri.host);
|
||||
if (!app) {
|
||||
stats.increment('sbc.terminations', ['sipStatus:404']);
|
||||
return res.send(404, {headers: {'X-Reason': 'no configured application'}});
|
||||
res.send(404, {headers: {'X-Reason': 'no configured application'}});
|
||||
return req.srf.endSession(req);
|
||||
}
|
||||
|
||||
req.locals = {
|
||||
@@ -171,7 +172,8 @@ module.exports = function(srf, logger) {
|
||||
const account = await lookupAccountBySipRealm(uri.host);
|
||||
if (!account) {
|
||||
stats.increment('sbc.terminations', ['sipStatus:404']);
|
||||
return res.send(404);
|
||||
res.send(404);
|
||||
return req.srf.endSession(req);
|
||||
}
|
||||
|
||||
/* if this is a dedicated SBC (static IP) only take calls for that account's sip realm */
|
||||
@@ -180,7 +182,8 @@ 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;
|
||||
return res.send(404);
|
||||
res.send(404);
|
||||
return req.srf.endSession(req);
|
||||
}
|
||||
req.locals = {
|
||||
account_sid: account.account_sid,
|
||||
@@ -261,13 +264,15 @@ module.exports = function(srf, logger) {
|
||||
account_sid,
|
||||
count: limit_sessions
|
||||
}).catch((err) => logger.info({err}, 'checkLimits: error writing alert'));
|
||||
return res.send(503, 'Maximum Calls In Progress');
|
||||
res.send(503, 'Maximum Calls In Progress');
|
||||
return req.srf.endSession(req);
|
||||
}
|
||||
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);
|
||||
}
|
||||
};
|
||||
|
||||
@@ -280,6 +285,7 @@ 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);
|
||||
}
|
||||
};
|
||||
|
||||
|
||||
Generated
+902
-358
File diff suppressed because it is too large
Load Diff
+13
-12
@@ -25,23 +25,24 @@
|
||||
"jslint": "eslint app.js lib"
|
||||
},
|
||||
"dependencies": {
|
||||
"@jambonz/db-helpers": "^0.6.12",
|
||||
"@jambonz/db-helpers": "^0.6.16",
|
||||
"@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",
|
||||
"@jambonz/http-health-check": "^0.0.1",
|
||||
"@jambonz/realtimedb-helpers": "^0.4.9",
|
||||
"@jambonz/rtpengine-utils": "^0.2.2",
|
||||
"@jambonz/stats-collector": "^0.1.6",
|
||||
"@jambonz/time-series": "^0.1.6",
|
||||
"aws-sdk": "^2.1036.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.1",
|
||||
"debug": "^4.3.3",
|
||||
"drachtio-fn-b2b-sugar": "0.0.12",
|
||||
"drachtio-srf": "^4.4.59",
|
||||
"pino": "^6.8.0",
|
||||
"rtpengine-client": "^0.2.0"
|
||||
"express": "^4.17.1",
|
||||
"pino": "^7.4.1",
|
||||
"rtpengine-client": "^0.2.0",
|
||||
"verify-aws-sns-signature": "^0.0.6",
|
||||
"xml2js": "^0.4.23"
|
||||
},
|
||||
"devDependencies": {
|
||||
"clear-module": "^4.1.1",
|
||||
|
||||
@@ -10,6 +10,7 @@ networks:
|
||||
services:
|
||||
mysql:
|
||||
image: mysql:5.7
|
||||
platform: linux/x86_64
|
||||
ports:
|
||||
- "3306:3306"
|
||||
environment:
|
||||
@@ -71,7 +72,7 @@ services:
|
||||
ipv4_address: 172.38.0.14
|
||||
|
||||
influxdb:
|
||||
image: influxdb:1.8-alpine
|
||||
image: influxdb:1.8
|
||||
ports:
|
||||
- "8086:8086"
|
||||
networks:
|
||||
|
||||
Reference in New Issue
Block a user