mirror of
https://github.com/jambonz/sbc-inbound.git
synced 2026-10-04 02:04:22 +00:00
Compare commits
110
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
91a29d2b2c | ||
|
|
27686a80e8 | ||
|
|
a9e79b23bd | ||
|
|
faafaa8549 | ||
|
|
b24ecf84d3 | ||
|
|
3360ddb3e6 | ||
|
|
b2d48c8954 | ||
|
|
1f1a4d2330 | ||
|
|
c34da3cfa8 | ||
|
|
5b88065925 | ||
|
|
438924ca36 | ||
|
|
c4d4d7bc0a | ||
|
|
ebcd0ce4a2 | ||
|
|
cf239ca7fa | ||
|
|
b072524585 | ||
|
|
7fb9966d20 | ||
|
|
42b3317206 | ||
|
|
db66d98d0a | ||
|
|
26169517dc | ||
|
|
c78ec892b9 | ||
|
|
6d237d43b3 | ||
|
|
44d298a4f3 | ||
|
|
030059a596 | ||
|
|
78b60525e2 | ||
|
|
dc9103cfa1 | ||
|
|
4dd4247dc2 | ||
|
|
055350903a | ||
|
|
7e83791ef0 | ||
|
|
dbedd7419e | ||
|
|
c4d07b517e | ||
|
|
2693077871 | ||
|
|
ad0912d302 | ||
|
|
e715433534 | ||
|
|
f315b1b41f | ||
|
|
bd09703732 | ||
|
|
363eb676a3 | ||
|
|
7b085e5763 | ||
|
|
04554fd60d | ||
|
|
9f1109c738 | ||
|
|
d27847747d | ||
|
|
88bdeffe4b | ||
|
|
3245fae069 | ||
|
|
213e84f59c | ||
|
|
ca0c9c157c | ||
|
|
996519404e | ||
|
|
acead419d5 | ||
|
|
911d208e0f | ||
|
|
c824367086 | ||
|
|
f39df615f2 | ||
|
|
58e80afc92 | ||
|
|
738f151066 | ||
|
|
a94f25b0bd | ||
|
|
a1542b161b | ||
|
|
4e961491a6 | ||
|
|
4e5f7ae908 | ||
|
|
798e070127 | ||
|
|
2d350d4850 | ||
|
|
baad125924 | ||
|
|
73528f5ce2 | ||
|
|
c8329a94f3 | ||
|
|
dceeef5549 | ||
|
|
b5551fffba | ||
|
|
d2b5597571 | ||
|
|
c9401ab3c8 | ||
|
|
cc5c712a5b | ||
|
|
a923227e4a | ||
|
|
5f7df4d135 | ||
|
|
7dd4d4a045 | ||
|
|
7382bf6bd6 | ||
|
|
abfce38150 | ||
|
|
d8035c978e | ||
|
|
91cd677ad8 | ||
|
|
17b53438e2 | ||
|
|
f1e2f192b7 | ||
|
|
0af8bcc348 | ||
|
|
e067125974 | ||
|
|
2dadde64f4 | ||
|
|
d5a1337811 | ||
|
|
a6571134cd | ||
|
|
e0e5c75496 | ||
|
|
c08e35c261 | ||
|
|
883c63723c | ||
|
|
13f78bf8d9 | ||
|
|
55056c1771 | ||
|
|
87ec5f8e09 | ||
|
|
317280befc | ||
|
|
7b97a0a137 | ||
|
|
75d8381ebb | ||
|
|
7750fdc3f2 | ||
|
|
fcf90916e1 | ||
|
|
4a0aa79456 | ||
|
|
1cdaae7773 | ||
|
|
82ac86389d | ||
|
|
e9184ad208 | ||
|
|
804bb890b5 | ||
|
|
a892a87eb5 | ||
|
|
56205dc852 | ||
|
|
927b7be637 | ||
|
|
63b482e562 | ||
|
|
a7d047a7e8 | ||
|
|
ea4f5ea0a8 | ||
|
|
13d8ee8c3c | ||
|
|
c717470c5e | ||
|
|
f16598e144 | ||
|
|
8584020d4c | ||
|
|
0b1b49bd59 | ||
|
|
e22b2ae63f | ||
|
|
2a5200673d | ||
|
|
13ca308df7 | ||
|
|
981788dc58 |
+1
-1
@@ -8,7 +8,7 @@
|
||||
"jsx": false,
|
||||
"modules": false
|
||||
},
|
||||
"ecmaVersion": 2018
|
||||
"ecmaVersion": 2020
|
||||
},
|
||||
"plugins": ["promise"],
|
||||
"rules": {
|
||||
|
||||
@@ -0,0 +1,51 @@
|
||||
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: 12
|
||||
- run: npm install
|
||||
node-version: 18.x
|
||||
- run: npm ci
|
||||
- run: npm run jslint
|
||||
- run: npm test
|
||||
|
||||
+1
-1
@@ -1,7 +1,7 @@
|
||||
# Logs
|
||||
logs
|
||||
*.log
|
||||
|
||||
.vscode/
|
||||
# Runtime data
|
||||
pids
|
||||
*.pid
|
||||
|
||||
Executable
+4
@@ -0,0 +1,4 @@
|
||||
#!/bin/sh
|
||||
. "$(dirname "$0")/_/husky.sh"
|
||||
|
||||
npm run jslint
|
||||
+18
-11
@@ -1,16 +1,23 @@
|
||||
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 --platform=linux/amd64 node:18.12.1-alpine3.16 as base
|
||||
|
||||
FROM node:alpine as app
|
||||
WORKDIR /opt/app
|
||||
COPY . /opt/app
|
||||
COPY --from=builder /opt/app/node_modules ./node_modules
|
||||
RUN apk --update --no-cache add --virtual .builds-deps build-base python3
|
||||
|
||||
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/
|
||||
|
||||
ARG NODE_ENV
|
||||
|
||||
ENV NODE_ENV $NODE_ENV
|
||||
|
||||
CMD [ "npm", "start" ]
|
||||
CMD [ "node", "app.js" ]
|
||||
|
||||
@@ -6,27 +6,30 @@ 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'});
|
||||
const logger = require('pino')(opts);
|
||||
const {
|
||||
writeCallCount,
|
||||
writeCallCountSP,
|
||||
writeCallCountApp,
|
||||
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 = [];
|
||||
@@ -34,6 +37,7 @@ const setName = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-sip`;
|
||||
|
||||
const {
|
||||
pool,
|
||||
ping,
|
||||
lookupAuthHook,
|
||||
lookupSipGatewayBySignalingAddress,
|
||||
addSbcAddress,
|
||||
@@ -41,7 +45,8 @@ const {
|
||||
lookupAppByTeamsTenant,
|
||||
lookupAccountBySipRealm,
|
||||
lookupAccountBySid,
|
||||
lookupAccountCapacitiesBySid
|
||||
lookupAccountCapacitiesBySid,
|
||||
queryCallLimits
|
||||
} = require('@jambonz/db-helpers')({
|
||||
host: process.env.JAMBONES_MYSQL_HOST,
|
||||
user: process.env.JAMBONES_MYSQL_USER,
|
||||
@@ -49,14 +54,30 @@ const {
|
||||
database: process.env.JAMBONES_MYSQL_DATABASE,
|
||||
connectionLimit: process.env.JAMBONES_MYSQL_CONNECTION_LIMIT || 10
|
||||
}, logger);
|
||||
const {createSet, retrieveSet, addToSet, removeFromSet, incrKey, decrKey} = require('@jambonz/realtimedb-helpers')({
|
||||
const {
|
||||
client: redisClient,
|
||||
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});
|
||||
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,
|
||||
dtmfListenPort: process.env.DTMF_LISTEN_PORT || 22224,
|
||||
protocol: ngProtocol
|
||||
});
|
||||
srf.locals = {...srf.locals,
|
||||
stats,
|
||||
writeCallCount,
|
||||
writeCallCountSP,
|
||||
writeCallCountApp,
|
||||
queryCdrs,
|
||||
writeCdrs,
|
||||
writeAlerts,
|
||||
@@ -65,13 +86,15 @@ srf.locals = {...srf.locals,
|
||||
getRtpEngine,
|
||||
dbHelpers: {
|
||||
pool,
|
||||
ping,
|
||||
lookupAuthHook,
|
||||
lookupSipGatewayBySignalingAddress,
|
||||
lookupAccountByPhoneNumber,
|
||||
lookupAppByTeamsTenant,
|
||||
lookupAccountBySid,
|
||||
lookupAccountBySipRealm,
|
||||
lookupAccountCapacitiesBySid
|
||||
lookupAccountCapacitiesBySid,
|
||||
queryCallLimits
|
||||
},
|
||||
realtimeDbHelpers: {
|
||||
createSet,
|
||||
@@ -80,23 +103,42 @@ srf.locals = {...srf.locals,
|
||||
retrieveSet
|
||||
}
|
||||
};
|
||||
srf.locals.getFeatureServer = require('./lib/fs-tracking')(srf, logger);
|
||||
const {
|
||||
getSPForAccount,
|
||||
wasOriginatedFromCarrier,
|
||||
getApplicationForDidAndCarrier,
|
||||
getOutboundGatewayForRefer
|
||||
} = require('./lib/db-utils')(srf, logger);
|
||||
srf.locals = {
|
||||
...srf.locals,
|
||||
getSPForAccount,
|
||||
wasOriginatedFromCarrier,
|
||||
getApplicationForDidAndCarrier,
|
||||
getOutboundGatewayForRefer,
|
||||
getFeatureServer: require('./lib/fs-tracking')(srf, logger)
|
||||
};
|
||||
const activeCallIds = srf.locals.activeCallIds;
|
||||
|
||||
const {
|
||||
initLocals,
|
||||
handleSipRec,
|
||||
identifyAccount,
|
||||
checkLimits,
|
||||
challengeDeviceCalls
|
||||
} = 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');
|
||||
logger.info(`connected to drachtio listening on ${hp}`);
|
||||
if (process.env.SBC_ACCOUNT_SID) return;
|
||||
|
||||
const hostports = hp.split(',');
|
||||
for (const hp of hostports) {
|
||||
@@ -104,11 +146,12 @@ if (process.env.DRACHTIO_HOST) {
|
||||
if (arr && 'udp' === arr[1] && !matcher.contains(arr[2])) {
|
||||
logger.info(`adding sbc public address to database: ${arr[2]}`);
|
||||
srf.locals.sipAddress = arr[2];
|
||||
addSbcAddress(arr[2]);
|
||||
if (!process.env.SBC_ACCOUNT_SID) 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();
|
||||
@@ -117,6 +160,9 @@ if (process.env.DRACHTIO_HOST) {
|
||||
});
|
||||
}
|
||||
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') {
|
||||
@@ -126,7 +172,13 @@ if (process.env.NODE_ENV === 'test') {
|
||||
}
|
||||
|
||||
/* install middleware */
|
||||
srf.use('invite', [initLocals, identifyAccount, checkLimits, challengeDeviceCalls]);
|
||||
srf.use('invite', [
|
||||
initLocals,
|
||||
handleSipRec,
|
||||
identifyAccount,
|
||||
checkLimits,
|
||||
challengeDeviceCalls
|
||||
]);
|
||||
|
||||
srf.invite((req, res) => {
|
||||
if (req.has('Replaces')) {
|
||||
@@ -149,35 +201,79 @@ srf.use((req, res, next, err) => {
|
||||
res.send(500);
|
||||
});
|
||||
|
||||
/* update call stats periodically */
|
||||
setInterval(() => {
|
||||
stats.gauge('sbc.sip.calls.count', activeCallIds.size, ['direction:inbound']);
|
||||
}, 20000);
|
||||
if (process.env.K8S || process.env.HTTP_PORT) {
|
||||
const PORT = process.env.HTTP_PORT || 3000;
|
||||
const healthCheck = require('@jambonz/http-health-check');
|
||||
|
||||
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;
|
||||
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) => {
|
||||
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}:${ngPort}`));
|
||||
}
|
||||
});
|
||||
};
|
||||
|
||||
/* update rtpengines periodically */
|
||||
if (process.env.JAMBONES_RTPENGINES) {
|
||||
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 */
|
||||
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 (!arrayCompare(newArray, rtpServers)) {
|
||||
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);
|
||||
}
|
||||
@@ -185,11 +281,11 @@ else {
|
||||
logger.error({err}, 'Error setting new rtpengines');
|
||||
}
|
||||
};
|
||||
|
||||
setInterval(() => {
|
||||
getActiveRtpServers();
|
||||
}, 30000);
|
||||
getActiveRtpServers();
|
||||
|
||||
}
|
||||
|
||||
const {lifecycleEmitter} = require('./lib/autoscale-manager')(logger);
|
||||
@@ -204,4 +300,12 @@ 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};
|
||||
|
||||
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);
|
||||
}
|
||||
})();
|
||||
@@ -3,6 +3,6 @@
|
||||
"DTLS": "off",
|
||||
"SDES": "off",
|
||||
"ICE": "remove",
|
||||
"flags": ["media handover"],
|
||||
"flags": ["media handover", "port latching"],
|
||||
"rtcp-mux": ["accept"]
|
||||
}
|
||||
@@ -1,16 +1,16 @@
|
||||
{
|
||||
"default": {
|
||||
"transport-protocol": "UDP/TLS/RTP/SAVPF",
|
||||
"ICE": "force",
|
||||
"ICE": "default",
|
||||
"SDES": "off",
|
||||
"flags": ["generate mid", "SDES-no", "media handover"],
|
||||
"flags": ["generate mid", "SDES-no", "media handover", "port latching"],
|
||||
"rtcp-mux": ["require"]
|
||||
},
|
||||
"teams": {
|
||||
"transport-protocol": "RTP/SAVP",
|
||||
"ICE": "force",
|
||||
"ICE": "default",
|
||||
"SDES": "off",
|
||||
"flags": ["generate mid", "SDES-no", "media handover"],
|
||||
"flags": ["generate mid", "SDES-no", "media handover", "port latching"],
|
||||
"rtcp-mux": ["accept"]
|
||||
}
|
||||
}
|
||||
@@ -1,7 +1,7 @@
|
||||
const Emitter = require('events');
|
||||
const bent = require('bent');
|
||||
const assert = require('assert');
|
||||
const PORT = process.env.AWS_SNS_PORT || 3001;
|
||||
const PORT = process.env.AWS_SNS_PORT || 3010;
|
||||
const {LifeCycleEvents} = require('./constants');
|
||||
const express = require('express');
|
||||
const app = express();
|
||||
@@ -21,6 +21,26 @@ class SnsNotifier extends Emitter {
|
||||
|
||||
this.logger = logger;
|
||||
}
|
||||
_doListen(logger, app, port, resolve) {
|
||||
return app.listen(port, () => {
|
||||
this.snsEndpoint = `http://${this.publicIp}:${port}`;
|
||||
logger.info(`SNS lifecycle server listening on http://localhost:${port}`);
|
||||
resolve(app);
|
||||
});
|
||||
}
|
||||
|
||||
_handleErrors(logger, app, resolve, reject, e) {
|
||||
if (e.code === 'EADDRINUSE' &&
|
||||
process.env.AWS_SNS_PORT_MAX &&
|
||||
e.port < process.env.AWS_SNS_PORT_MAX) {
|
||||
|
||||
logger.info(`SNS lifecycle server failed to bind port on ${e.port}, will try next port`);
|
||||
const server = this._doListen(logger, app, ++e.port, resolve);
|
||||
server.on('error', this._handleErrors.bind(this, logger, app, resolve, reject));
|
||||
return;
|
||||
}
|
||||
reject(e);
|
||||
}
|
||||
|
||||
async _handlePost(req, res) {
|
||||
try {
|
||||
@@ -45,6 +65,7 @@ class SnsNotifier extends Emitter {
|
||||
}, 'response from SNS SubscribeURL');
|
||||
const data = await this.describeInstance();
|
||||
this.lifecycleState = data.AutoScalingInstances[0].LifecycleState;
|
||||
this.emit('SubscriptionConfirmation', {publicIp: this.publicIp});
|
||||
break;
|
||||
|
||||
case 'Notification':
|
||||
@@ -80,14 +101,12 @@ class SnsNotifier extends Emitter {
|
||||
|
||||
async init() {
|
||||
try {
|
||||
this.logger.info('SnsNotifier: retrieving instance data');
|
||||
this.logger.debug('SnsNotifier: retrieving instance data');
|
||||
this.instanceId = await getString('http://169.254.169.254/latest/meta-data/instance-id');
|
||||
this.publicIp = await getString('http://169.254.169.254/latest/meta-data/public-ipv4');
|
||||
this.snsEndpoint = `http://${this.publicIp}:${PORT}`;
|
||||
this.logger.info({
|
||||
instanceId: this.instanceId,
|
||||
publicIp: this.publicIp,
|
||||
snsEndpoint: this.snsEndpoint
|
||||
publicIp: this.publicIp
|
||||
}, 'retrieved AWS instance data');
|
||||
|
||||
// start listening
|
||||
@@ -99,7 +118,10 @@ class SnsNotifier extends Emitter {
|
||||
this.logger.error(err, 'burped error');
|
||||
res.status(err.status || 500).json({msg: err.message});
|
||||
});
|
||||
app.listen(PORT);
|
||||
return new Promise((resolve, reject) => {
|
||||
const server = this._doListen(this.logger, app, PORT, resolve);
|
||||
server.on('error', this._handleErrors.bind(this, this.logger, app, resolve, reject));
|
||||
});
|
||||
|
||||
} catch (err) {
|
||||
this.logger.error({err}, 'Error retrieving AWS instance metadata');
|
||||
|
||||
+526
-45
@@ -1,7 +1,16 @@
|
||||
const Emitter = require('events');
|
||||
const {makeRtpEngineOpts, SdpWantsSrtp, makeCallCountKey} = require('./utils');
|
||||
const SrsClient = require('@jambonz/siprec-client-utils');
|
||||
const {
|
||||
makeRtpEngineOpts,
|
||||
SdpWantsSrtp,
|
||||
SdpWantsSDES,
|
||||
nudgeCallCounts,
|
||||
roundTripTime,
|
||||
parseConnectionIp
|
||||
} = require('./utils');
|
||||
|
||||
const {forwardInDialogRequests} = require('drachtio-fn-b2b-sugar');
|
||||
const {parseUri, SipError} = require('drachtio-srf');
|
||||
const {parseUri, stringifyUri, SipError} = require('drachtio-srf');
|
||||
const debug = require('debug')('jambonz:sbc-inbound');
|
||||
const MS_TEAMS_USER_AGENT = 'Microsoft.PSTNHub.SIPProxy';
|
||||
const MS_TEAMS_SIP_ENDPOINT = 'sip.pstnhub.microsoft.com';
|
||||
@@ -17,6 +26,20 @@ const createBLegFromHeader = (req) => {
|
||||
return '<sip:anonymous@localhost>';
|
||||
};
|
||||
|
||||
const createSiprecBody = (headers, sdp, type, content) => {
|
||||
const sep = 'uniqueBoundary';
|
||||
headers['Content-Type'] = `multipart/mixed;boundary="${sep}"`;
|
||||
return `--${sep}\r
|
||||
Content-Type: application/sdp\r
|
||||
\r
|
||||
${sdp}\r
|
||||
--${sep}\r
|
||||
Content-Type: ${type}\r
|
||||
Content-Disposition: recording-session\r
|
||||
\r
|
||||
${content}`;
|
||||
};
|
||||
|
||||
class CallSession extends Emitter {
|
||||
constructor(logger, req, res) {
|
||||
super();
|
||||
@@ -24,6 +47,8 @@ class CallSession extends Emitter {
|
||||
this.res = res;
|
||||
this.srf = req.srf;
|
||||
this.logger = logger.child({callId: req.get('Call-ID')});
|
||||
this.siprec = req.locals.siprec;
|
||||
this.xml = req.locals.xml;
|
||||
|
||||
this.getRtpEngine = req.srf.locals.getRtpEngine;
|
||||
this.getFeatureServer = req.srf.locals.getFeatureServer;
|
||||
@@ -32,14 +57,54 @@ class CallSession extends Emitter {
|
||||
this.activeCallIds = this.srf.locals.activeCallIds;
|
||||
|
||||
this.decrKey = req.srf.locals.realtimeDbHelpers.decrKey;
|
||||
this.callCountKey = makeCallCountKey(req.locals.account_sid);
|
||||
this._mediaReleased = false;
|
||||
|
||||
this.application_sid = req.locals.application_sid;
|
||||
this.account_sid = req.locals.account_sid;
|
||||
this.service_provider_sid = req.locals.service_provider_sid;
|
||||
}
|
||||
|
||||
get isFromMSTeams() {
|
||||
return !!this.req.locals.msTeamsTenantFqdn;
|
||||
}
|
||||
|
||||
get privateSipAddress() {
|
||||
return this.srf.locals.privateSipAddress;
|
||||
}
|
||||
|
||||
get isMediaReleased() {
|
||||
return this._mediaReleased;
|
||||
}
|
||||
|
||||
get callerIsUsingSrtp() {
|
||||
const tp = this.rtpEngineOpts?.uas?.mediaOpts['transport-protocol'];
|
||||
return tp && -1 !== tp.indexOf('SAVP');
|
||||
}
|
||||
|
||||
get isFive9VoiceStream() {
|
||||
return this.req.has('X-Five9-StreamingPairId');
|
||||
}
|
||||
|
||||
get isPossibleWebRtcClient() {
|
||||
return this.req.locals.isPossibleWebRtcClient;
|
||||
}
|
||||
|
||||
subscribeForDTMF(dlg) {
|
||||
if (!this._subscribedForDTMF) {
|
||||
this._subscribedForDTMF = true;
|
||||
this.subscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag,
|
||||
this._onDTMF.bind(this, dlg));
|
||||
}
|
||||
}
|
||||
unsubscribeForDTMF() {
|
||||
if (this._subscribedForDTMF) {
|
||||
this._subscribedForDTMF = false;
|
||||
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
|
||||
}
|
||||
}
|
||||
|
||||
async connect() {
|
||||
const {sdp} = this.req.locals;
|
||||
this.logger.info('inbound call accepted for routing');
|
||||
const engine = this.getRtpEngine();
|
||||
if (!engine) {
|
||||
@@ -49,10 +114,34 @@ class CallSession extends Emitter {
|
||||
return this.res.send(480);
|
||||
}
|
||||
debug(`got engine: ${JSON.stringify(engine)}`);
|
||||
const {offer, answer, del} = engine;
|
||||
const {
|
||||
offer,
|
||||
answer,
|
||||
del,
|
||||
blockMedia,
|
||||
unblockMedia,
|
||||
blockDTMF,
|
||||
unblockDTMF,
|
||||
playDTMF,
|
||||
subscribeDTMF,
|
||||
unsubscribeDTMF,
|
||||
subscribeRequest,
|
||||
subscribeAnswer,
|
||||
unsubscribe
|
||||
} = engine;
|
||||
this.offer = offer;
|
||||
this.answer = answer;
|
||||
this.del = del;
|
||||
this.blockMedia = blockMedia;
|
||||
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;
|
||||
|
||||
const featureServer = await this.getFeatureServer();
|
||||
if (!featureServer) {
|
||||
@@ -63,7 +152,9 @@ class CallSession extends Emitter {
|
||||
}
|
||||
this.logger.debug(`using feature server ${featureServer}`);
|
||||
|
||||
this.rtpEngineOpts = makeRtpEngineOpts(this.req, SdpWantsSrtp(this.req.body), false, this.isFromMSTeams);
|
||||
const wantsSrtp = this.req.locals.possibleWebRtcClient = SdpWantsSrtp(sdp);
|
||||
const wantsSDES = SdpWantsSDES(sdp);
|
||||
this.rtpEngineOpts = makeRtpEngineOpts(this.req, wantsSrtp, false, this.isFromMSTeams || wantsSDES);
|
||||
this.rtpEngineResource = {destroy: this.del.bind(null, this.rtpEngineOpts.common)};
|
||||
const obj = parseUri(this.req.uri);
|
||||
let proxy, host, uri;
|
||||
@@ -88,25 +179,40 @@ class CallSession extends Emitter {
|
||||
...this.rtpEngineOpts.uac.mediaOpts,
|
||||
'from-tag': this.rtpEngineOpts.uas.tag,
|
||||
direction: ['public', 'private'],
|
||||
sdp: this.req.body
|
||||
sdp
|
||||
};
|
||||
const startAt = process.hrtime();
|
||||
const response = await this.offer(opts);
|
||||
this.logger.debug({opts, response}, 'response from rtpengine to offer');
|
||||
this.rtpengineIp = opts.sdp ? parseConnectionIp(opts.sdp) : 'undefined';
|
||||
const rtt = roundTripTime(startAt);
|
||||
this.stats.histogram('app.rtpengine.response_time', rtt, [
|
||||
'direction:inbound', 'command:offer', `rtpengine:${this.rtpengineIp}`]);
|
||||
this.logger.debug({opts, response, rtt, rtpengine: this.rtpengineIp}, 'response from rtpengine to offer');
|
||||
if ('ok' !== response.result) {
|
||||
this.logger.error({}, `rtpengine offer failed with ${JSON.stringify(response)}`);
|
||||
throw new Error('rtpengine failed: answer');
|
||||
}
|
||||
|
||||
// now send the INVITE in towards the feature servers
|
||||
const headers = {
|
||||
let 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}:${this.req.source_port}`
|
||||
'X-Forwarded-For': `${this.req.source_address}`
|
||||
};
|
||||
if (this.privateSipAddress) headers = {...headers, Contact: `<sip:${this.privateSipAddress}>`};
|
||||
|
||||
const spdOfferB = this.siprec && this.xml ?
|
||||
createSiprecBody(headers, response.sdp, this.xml.type, this.xml.content) :
|
||||
response.sdp;
|
||||
|
||||
const responseHeaders = {};
|
||||
if (this.req.locals.carrier) Object.assign(headers, {'X-Originating-Carrier': this.req.locals.carrier});
|
||||
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.msTeamsTenantFqdn) {
|
||||
Object.assign(headers, {'X-MS-Teams-Tenant-FQDN': this.req.locals.msTeamsTenantFqdn});
|
||||
|
||||
@@ -131,14 +237,22 @@ class CallSession extends Emitter {
|
||||
|
||||
if (this.req.canceled) throw new Error('call canceled');
|
||||
|
||||
// now send the INVITE in towards the feature servers
|
||||
debug(`sending INVITE to ${proxy} with ${uri}`);
|
||||
const {uas, uac} = await this.srf.createB2BUA(this.req, this.res, uri, {
|
||||
proxy,
|
||||
headers,
|
||||
responseHeaders,
|
||||
proxyRequestHeaders: ['all', '-Authorization', '-Max-Forwards', '-Record-Route', '-Session-Expires', 'Min-SE'],
|
||||
proxyResponseHeaders: ['all'],
|
||||
localSdpB: response.sdp,
|
||||
proxyRequestHeaders: [
|
||||
'all',
|
||||
'-Authorization',
|
||||
'-Max-Forwards',
|
||||
'-Record-Route',
|
||||
'-Session-Expires',
|
||||
'-X-Subspace-Forwarded-For'
|
||||
],
|
||||
proxyResponseHeaders: ['all', '-X-Trace-ID'],
|
||||
localSdpB: spdOfferB,
|
||||
localSdpA: async(sdp, res) => {
|
||||
this.rtpEngineOpts.uac.tag = res.getParsedHeader('To').params.tag;
|
||||
const opts = {
|
||||
@@ -148,11 +262,27 @@ class CallSession extends Emitter {
|
||||
'to-tag': this.rtpEngineOpts.uac.tag,
|
||||
sdp
|
||||
};
|
||||
const startAt = process.hrtime();
|
||||
const response = await this.answer(opts);
|
||||
this.logger.debug({response, opts}, 'response from rtpengine to answer');
|
||||
const rtt = roundTripTime(startAt);
|
||||
this.stats.histogram('app.rtpengine.response_time', rtt, [
|
||||
'direction:inbound', 'command:answer', `rtpengine:${this.rtpengineIp}`]);
|
||||
if ('ok' !== response.result) {
|
||||
this.logger.error(`rtpengine answer failed with ${JSON.stringify(response)}`);
|
||||
throw new Error('rtpengine failed: answer');
|
||||
}
|
||||
/* special case: Five9 Voicestream calls do not advertise a:sendonly, though they should */
|
||||
if (this.isFive9VoiceStream) {
|
||||
const opts = {
|
||||
...this.rtpEngineOpts.common,
|
||||
'from-tag':this.rtpEngineOpts.uac.tag
|
||||
};
|
||||
this.logger.info('Voicestream call from Five9, blocking audio in the reverse direction');
|
||||
const response = await Promise.all([this.blockMedia(opts), this.blockDTMF(opts)]);
|
||||
this.logger.debug({response}, 'response to blockMedia/blockDTMF');
|
||||
}
|
||||
|
||||
return response.sdp;
|
||||
}
|
||||
});
|
||||
@@ -163,29 +293,39 @@ class CallSession extends Emitter {
|
||||
this._setHandlers({uas, uac});
|
||||
return;
|
||||
} catch (err) {
|
||||
this.rtpEngineResource.destroy();
|
||||
this.rtpEngineResource.destroy().catch((err) => this.logger.info({err}, 'Error destroying rtpe after failure'));
|
||||
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}`);
|
||||
return this.emit('failed');
|
||||
this.emit('failed');
|
||||
}
|
||||
else if (err.message !== 'call canceled') {
|
||||
this.logger.error(err, 'unexpected error routing inbound call');
|
||||
}
|
||||
this.srf.endSession(this.req);
|
||||
}
|
||||
}
|
||||
|
||||
_setDlgHandlers(dlg) {
|
||||
this.activeCallIds.set(this.req.get('Call-ID'), this);
|
||||
const {callId} = dlg.sip;
|
||||
this.activeCallIds.set(callId, this);
|
||||
if (this.isPossibleWebRtcClient) this.subscribeForDTMF(dlg);
|
||||
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(this.req.get('Call-ID'));
|
||||
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
|
||||
@@ -198,27 +338,43 @@ class CallSession extends Emitter {
|
||||
const tags = ['accepted:yes', 'sipStatus:200', `originator:${this.req.locals.originator}`];
|
||||
this.stats.increment('sbc.terminations', tags);
|
||||
this.activeCallIds.set(this.req.get('Call-ID'), this);
|
||||
const call_sid = uac.res?.get('X-Call-Sid');
|
||||
const application_sid = this.application_sid || uac.res?.get('X-Application-Sid');
|
||||
if (this.req.locals.cdr) {
|
||||
this.req.locals.cdr = {
|
||||
...this.req.locals.cdr,
|
||||
answered: true,
|
||||
answered_at: callStart
|
||||
answered_at: callStart,
|
||||
...(call_sid && {call_sid}),
|
||||
...(application_sid && {application_sid}),
|
||||
trace_id: uac.res?.get('X-Trace-ID') || '00000000000000000000000000000000'
|
||||
};
|
||||
}
|
||||
this.uas = uas;
|
||||
this.uac = uac;
|
||||
[uas, uac].forEach((dlg) => {
|
||||
dlg.on('destroy', () => {
|
||||
this.logger.info('call ended with normal termination');
|
||||
dlg.on('destroy', async() => {
|
||||
const other = dlg.other;
|
||||
this.rtpEngineResource.destroy().catch((err) => {});
|
||||
this.activeCallIds.delete(this.req.get('Call-ID'));
|
||||
dlg.other.destroy().catch((e) => {});
|
||||
try {
|
||||
await other.destroy();
|
||||
} catch (err) {}
|
||||
this.unsubscribeForDTMF();
|
||||
|
||||
if (process.env.JAMBONES_HOSTING) {
|
||||
this.decrKey(this.callCountKey)
|
||||
.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'));
|
||||
|
||||
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 */
|
||||
@@ -227,24 +383,81 @@ class CallSession extends Emitter {
|
||||
const trunk = ['trunk', 'teams'].includes(this.req.locals.originator) ?
|
||||
this.req.locals.carrier :
|
||||
this.req.locals.originator;
|
||||
const cdr = {...this.req.locals.cdr,
|
||||
terminated_at: now,
|
||||
termination_reason: dlg.type === 'uas' ? 'caller hungup' : 'called party hungup',
|
||||
sip_status: 200,
|
||||
duration: Math.floor((now - callStart) / 1000),
|
||||
trunk
|
||||
};
|
||||
this.logger.info({cdr}, 'going to write a cdr now..');
|
||||
this.writeCdrs({...this.req.locals.cdr,
|
||||
terminated_at: now,
|
||||
termination_reason: dlg.type === 'uas' ? 'caller hungup' : 'called party hungup',
|
||||
sip_status: 200,
|
||||
duration: Math.floor((now - callStart) / 1000),
|
||||
trunk
|
||||
}).catch((err) => this.logger.error({err}, 'Error writing cdr for completed call'));
|
||||
})
|
||||
.then(() => this.logger.debug('successfully wrote cdr'))
|
||||
.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);
|
||||
});
|
||||
});
|
||||
|
||||
if (this.isPossibleWebRtcClient) this.subscribeForDTMF(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, ['info', 'notify', 'options', 'message']);
|
||||
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');
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -254,7 +467,19 @@ class CallSession extends Emitter {
|
||||
*/
|
||||
async replaces(req, res) {
|
||||
try {
|
||||
let opts = Object.assign(this.rtpEngineOpts.offer, {sdp: req.body});
|
||||
const fromTag = this.rtpEngineOpts.uas.tag;
|
||||
const toTag = this.rtpEngineOpts.uac.tag;
|
||||
const offerMedia = this.rtpEngineOpts.uac.mediaOpts;
|
||||
const answerMedia = this.rtpEngineOpts.uas.mediaOpts;
|
||||
const direction = ['public', 'private'];
|
||||
let opts = {
|
||||
...this.rtpEngineOpts.common,
|
||||
...offerMedia,
|
||||
'from-tag': fromTag,
|
||||
'to-tag': toTag,
|
||||
direction,
|
||||
sdp: req.body,
|
||||
};
|
||||
let response = await this.offer(opts);
|
||||
if ('ok' !== response.result) {
|
||||
res.send(488);
|
||||
@@ -262,8 +487,13 @@ class CallSession extends Emitter {
|
||||
}
|
||||
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 = {
|
||||
...this.rtpEngineOpts.common,
|
||||
...answerMedia,
|
||||
'from-tag': fromTag,
|
||||
'to-tag': toTag,
|
||||
sdp
|
||||
};
|
||||
response = await this.answer(opts);
|
||||
if ('ok' !== response.result) {
|
||||
res.send(488);
|
||||
@@ -281,6 +511,8 @@ class CallSession extends Emitter {
|
||||
});
|
||||
}
|
||||
|
||||
this.unsubscribeForDTMF();
|
||||
|
||||
const uas = await this.srf.createUAS(req, res, {
|
||||
localSdp: response.sdp,
|
||||
headers
|
||||
@@ -307,26 +539,49 @@ class CallSession extends Emitter {
|
||||
res.send(200, {body: dlg.local.sdp});
|
||||
return;
|
||||
}
|
||||
const offeredSdp = Array.isArray(req.payload) && req.payload.length > 1 ?
|
||||
req.payload.find((p) => p.type === 'application/sdp').content :
|
||||
req.body;
|
||||
|
||||
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' ? ['public', 'private'] : ['private', 'public'];
|
||||
if (isReleasingMedia) {
|
||||
if (!offerMedia.flags.includes('asymmetric')) offerMedia.flags.push('asymmetric');
|
||||
offerMedia.flags = offerMedia.flags.filter((f) => f !== 'media handover');
|
||||
}
|
||||
let opts = {
|
||||
...this.rtpEngineOpts.common,
|
||||
...offerMedia,
|
||||
'from-tag': fromTag,
|
||||
'to-tag': toTag,
|
||||
direction,
|
||||
sdp: req.body,
|
||||
sdp: offeredSdp,
|
||||
};
|
||||
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)}`);
|
||||
}
|
||||
const sdp = await dlg.other.modify(response.sdp);
|
||||
|
||||
/* if this is a re-invite from the FS to change media anchoring, avoid sending the reinvite out */
|
||||
let sdp;
|
||||
if (isReleasingMedia && !this.callerIsUsingSrtp) {
|
||||
this.logger.info({response}, `got a reinvite from FS to ${reason}`);
|
||||
sdp = dlg.other.remote.sdp;
|
||||
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);
|
||||
}
|
||||
opts = {
|
||||
...this.rtpEngineOpts.common,
|
||||
...answerMedia,
|
||||
@@ -345,34 +600,247 @@ class CallSession extends Emitter {
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
const contentType = req.get('Content-Type');
|
||||
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}`);
|
||||
|
||||
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,
|
||||
direction: 'inbound',
|
||||
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 if (dlg.type === 'uas' && ['application/dtmf-relay', 'application/dtmf'].includes(contentType)) {
|
||||
const arr = /Signal=\s*([0-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.uas.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 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) {
|
||||
if (this.srsClient) {
|
||||
this.srsClient = null;
|
||||
}
|
||||
res.send(500);
|
||||
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}, 'received REFER from feature server');
|
||||
this.logger.info({uri, referTo, headers: req.headers}, 'received REFER from feature server');
|
||||
const arr = /context-(.*)/.exec(uri.user);
|
||||
if (!arr) {
|
||||
this.logger.info(`invalid Refer-To header: ${referTo.uri}`);
|
||||
return res.send(501);
|
||||
/* 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);
|
||||
const leaveReferToAlone = req.has('X-Refer-To-Leave-Untouched');
|
||||
if (leaveReferToAlone) {
|
||||
this.logger.debug({referTo}, 'passing Refer-To header through untouched');
|
||||
}
|
||||
else {
|
||||
const isDotDecimal = /^(?:[0-9]{1,3}\.){3}[0-9]{1,3}$/.test(uri.host);
|
||||
let selectedGateway = false;
|
||||
let e164 = false;
|
||||
if (gateway && isDotDecimal) {
|
||||
/* 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 && isDotDecimal) {
|
||||
uri.host = this.req.source_address;
|
||||
uri.port = this.req.source_port;
|
||||
}
|
||||
if (e164 && !uri.user.startsWith('+')) {
|
||||
uri.user = `+${uri.user}`;
|
||||
}
|
||||
}
|
||||
// eslint-disable-next-line no-unused-vars
|
||||
const {via, from, to, 'call-id':callid, cseq, 'max-forwards':maxforwards,
|
||||
// eslint-disable-next-line no-unused-vars
|
||||
'content-length':contentlength, 'refer-to':_referto, 'referred-by':_referredby,
|
||||
// eslint-disable-next-line no-unused-vars
|
||||
'X-Refer-To-Leave-Untouched': _leave,
|
||||
...customHeaders
|
||||
} = req.headers;
|
||||
|
||||
const response = await this.uas.request({
|
||||
method: 'REFER',
|
||||
headers: {
|
||||
'Refer-To': stringifyUri(uri),
|
||||
'Referred-By': stringifyUri(u),
|
||||
...customHeaders
|
||||
}
|
||||
});
|
||||
return res.send(response.status);
|
||||
}
|
||||
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 uac = await this.srf.createUAC(referTo.uri, {localSdp: dlg.local.sdp, headers});
|
||||
this.uac = uac;
|
||||
uac.other = this.uas;
|
||||
this.uas.other = uac;
|
||||
uac.on('modify', this._onFeatureServerReinvite.bind(this, uac));
|
||||
uac.on('modify', this._onReinvite.bind(this, uac));
|
||||
uac.on('refer', this._onFeatureServerTransfer.bind(this, uac));
|
||||
uac.on('destroy', () => {
|
||||
this.logger.info('call ended with normal termination');
|
||||
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(() => {});
|
||||
@@ -383,7 +851,7 @@ class CallSession extends Emitter {
|
||||
const response = await this.answer(opts);
|
||||
if ('ok' !== response.result) {
|
||||
res.send(488);
|
||||
throw new Error(`_onFeatureServerReinvite: rtpengine failed: ${JSON.stringify(response)}`);
|
||||
throw new Error(`_onFeatureServerTransfer: rtpengine failed: ${JSON.stringify(response)}`);
|
||||
}
|
||||
this.logger.info('successfully moved call to new feature server');
|
||||
} catch (err) {
|
||||
@@ -444,6 +912,7 @@ class CallSession extends Emitter {
|
||||
|
||||
// successfully connected
|
||||
this.logger.info('successfully connected new call leg for REFER');
|
||||
this.unsubscribeForDTMF();
|
||||
this.referInvite = null;
|
||||
sendNotify(this.uas, '200 OK');
|
||||
this.uas.destroy();
|
||||
@@ -457,8 +926,20 @@ class CallSession extends Emitter {
|
||||
}
|
||||
}
|
||||
else {
|
||||
// TODO: forward on to feature server
|
||||
res.send(501);
|
||||
/* 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');
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+214
-45
@@ -3,6 +3,8 @@ const CIDRMatcher = require('cidr-matcher');
|
||||
const {parseUri} = require('drachtio-srf');
|
||||
const {normalizeDID} = require('./utils');
|
||||
|
||||
const sqlSelectSPForAccount = 'SELECT service_provider_sid FROM accounts WHERE account_sid = ?';
|
||||
|
||||
const sqlSelectAllCarriersForAccountByRealm =
|
||||
`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
|
||||
@@ -11,9 +13,18 @@ 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
|
||||
vc.account_sid, vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask
|
||||
FROM sip_gateways sg, voip_carriers vc
|
||||
WHERE sg.voip_carrier_sid = vc.voip_carrier_sid
|
||||
AND vc.service_provider_sid IS NOT NULL
|
||||
@@ -35,6 +46,27 @@ 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
|
||||
WHERE sg.voip_carrier_sid = ?
|
||||
AND sg.voip_carrier_sid = vc.voip_carrier_sid
|
||||
AND outbound = 1`;
|
||||
|
||||
const sqlSelectCarrierRequiringRegistration = `
|
||||
SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.service_provider_sid, vc.account_sid,
|
||||
vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask
|
||||
FROM sip_gateways sg, voip_carriers vc
|
||||
WHERE sg.voip_carrier_sid = vc.voip_carrier_sid
|
||||
AND vc.requires_register = 1
|
||||
AND vc.is_active = 1
|
||||
AND vc.register_sip_realm = ?
|
||||
AND vc.register_username = ?`;
|
||||
|
||||
const gatewayMatchesSourceAddress = (source_address, gw) => {
|
||||
if (32 === gw.netmask && gw.ipv4 === source_address) return true;
|
||||
if (gw.netmask < 32) {
|
||||
@@ -48,13 +80,41 @@ module.exports = (srf, logger) => {
|
||||
const {pool} = srf.locals.dbHelpers;
|
||||
const pp = pool.promise();
|
||||
|
||||
const getSPForAccount = async(account_sid) => {
|
||||
const [r] = await pp.query(sqlSelectSPForAccount, [account_sid]);
|
||||
if (0 === r.length) return null;
|
||||
return r[0].service_provider_sid;
|
||||
};
|
||||
|
||||
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);
|
||||
|
||||
try {
|
||||
/* straight DID match */
|
||||
const [r] = await pp.query(sqlQueryApplicationByDid, [did, voip_carrier_sid]);
|
||||
if (0 === r.length) return null;
|
||||
return r[0].application_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;
|
||||
} catch (err) {
|
||||
logger.error({err}, 'getApplicationForDidAndCarrier');
|
||||
}
|
||||
@@ -65,66 +125,175 @@ module.exports = (srf, logger) => {
|
||||
const uri = parseUri(req.uri);
|
||||
const isDotDecimal = /^(?:[0-9]{1,3}\.){3}[0-9]{1,3}$/.test(uri.host);
|
||||
|
||||
if (isDotDecimal) {
|
||||
if (process.env.JAMBONES_HOSTING) {
|
||||
if (!process.env.SBC_ACCOUNT_SID) return failure;
|
||||
if (!isDotDecimal) {
|
||||
/**
|
||||
* The host part of the SIP URI is not a dot-decimal IP address,
|
||||
* so this can be one of two things:
|
||||
* (1) a sip realm value associate with an account, or
|
||||
* (2) a carrier name for a carrier that we send outbound registrations to
|
||||
*
|
||||
* Let's look for case #1 first...
|
||||
*/
|
||||
|
||||
/* look for carrier only within that account */
|
||||
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};
|
||||
/* 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 selected = gw.find(gatewayMatchesSourceAddress.bind(null, req.source_address));
|
||||
if (selected) {
|
||||
const [a] = await pp.query(sqlAccountByRealm, uri.host);
|
||||
if (0 === a.length) return failure;
|
||||
return {
|
||||
fromCarrier: true,
|
||||
gateway: selected,
|
||||
service_provider_sid: a[0].service_provider_sid,
|
||||
account_sid: a[0].account_sid,
|
||||
application_sid: selected.application_sid,
|
||||
account: a[0]
|
||||
};
|
||||
}
|
||||
else {
|
||||
/* we may have a carrier at the service provider level */
|
||||
const [gw] = await pp.query(sqlSelectAllGatewaysForSP);
|
||||
const matches = gw.filter(gatewayMatchesSourceAddress.bind(null, req.source_address));
|
||||
if (matches.length) {
|
||||
/* we have one or more carriers that match. Now we need to find one with a provisioned phone number */
|
||||
const vc_sids = matches.map((m) => `'${m.voip_carrier_sid}'`).join(',');
|
||||
const did = normalizeDID(req.calledNumber);
|
||||
const sql = `SELECT * FROM phone_numbers WHERE number = ${did} AND voip_carrier_sid IN (${vc_sids})`;
|
||||
logger.debug({matches, sql, did, vc_sids}, 'looking up DID');
|
||||
|
||||
const [r] = await pp.query(sql);
|
||||
if (0 === r.length) {
|
||||
/* 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);
|
||||
assert(accounts.length);
|
||||
/* no match, so let's look for case #2 */
|
||||
try {
|
||||
logger.info({
|
||||
host: uri.host,
|
||||
user: uri.user
|
||||
}, 'sip realm is not associated with an account, checking carriers');
|
||||
const [gw] = await pp.query(sqlSelectCarrierRequiringRegistration, [uri.host, uri.user]);
|
||||
const matches = gw.filter(gatewayMatchesSourceAddress.bind(null, req.source_address));
|
||||
if (1 === matches.length) {
|
||||
// bingo
|
||||
//TODO: this assumes the carrier is associate to an account, not an SP
|
||||
//if the carrier is associated with an SP (which would mean we
|
||||
//must see a dialed number in the To header, not the register username),
|
||||
//then we need to look up the account based on the dialed number in the To header
|
||||
const [a] = await pp.query(sqlAccountBySid, matches[0].account_sid);
|
||||
if (0 === a.length) return failure;
|
||||
logger.debug({matches}, `found registration carrier using ${uri.host} and ${uri.user}`);
|
||||
return {
|
||||
fromCarrier: true,
|
||||
gateway,
|
||||
account_sid: r[0].account_sid,
|
||||
application_sid: r[0].application_sid,
|
||||
account: accounts[0]
|
||||
gateway: matches[0],
|
||||
service_provider_sid: a[0].service_provider_sid,
|
||||
account_sid: a[0].account_sid,
|
||||
application_sid: matches[0].application_sid,
|
||||
account: a[0]
|
||||
};
|
||||
}
|
||||
return failure;
|
||||
else if (matches.length > 1) {
|
||||
logger.warn({matches, source_address: req.source_address}, 'multiple gateways match source address');
|
||||
return {
|
||||
fromCarrier: true,
|
||||
error: 'Multiple gateways match registration carrier source address'
|
||||
};
|
||||
}
|
||||
} catch (err) {
|
||||
logger.info({err, host: uri.host, user: uri.user}, 'Error looking up carrier by host and user');
|
||||
}
|
||||
/* no match, so fall through */
|
||||
}
|
||||
|
||||
/* get all the carriers and gateways for the account owning this sip realm */
|
||||
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);
|
||||
if (0 === a.length) return failure;
|
||||
if (isDotDecimal && process.env.JAMBONES_HOSTING) {
|
||||
if (!process.env.SBC_ACCOUNT_SID) return failure;
|
||||
|
||||
/* look for carrier only within that account */
|
||||
const [r] = await pp.query(sqlCarriersForAccountBySid,
|
||||
[process.env.SBC_ACCOUNT_SID, req.source_address, req.source_port]);
|
||||
if (0 === r.length) return failure;
|
||||
const service_provider_sid = await getSPForAccount(process.env.SBC_ACCOUNT_SID);
|
||||
return {
|
||||
fromCarrier: true,
|
||||
gateway: selected,
|
||||
account_sid: a[0].account_sid,
|
||||
application_sid: selected.application_sid,
|
||||
account: a[0]
|
||||
gateway: r[0],
|
||||
account_sid: process.env.SBC_ACCOUNT_SID,
|
||||
service_provider_sid
|
||||
};
|
||||
}
|
||||
else {
|
||||
/* find all carrier entries that have an inbound gateway matching the source IP */
|
||||
const [gw] = await pp.query(sqlSelectAllGatewaysForSP);
|
||||
const matches = gw.filter(gatewayMatchesSourceAddress.bind(null, req.source_address));
|
||||
if (matches.length) {
|
||||
/* we have one or more matches. Now check for one with a provisioned phone number matching the DID */
|
||||
const vc_sids = matches.map((m) => `'${m.voip_carrier_sid}'`).join(',');
|
||||
const did = normalizeDID(req.calledNumber);
|
||||
const sql = `SELECT * FROM phone_numbers WHERE number = '${did}' AND voip_carrier_sid IN (${vc_sids})`;
|
||||
logger.debug({matches, sql, did, vc_sids}, 'looking up DID');
|
||||
|
||||
const [r] = await pp.query(sql);
|
||||
if (0 === r.length) {
|
||||
/* came from a provisioned carrier, but the dialed number is not provisioned.
|
||||
check if we have an account with default routing of that carrier to an application
|
||||
*/
|
||||
const accountLevelGateways = matches.filter((m) => m.account_sid && m.application_sid);
|
||||
if (accountLevelGateways.length > 1) {
|
||||
logger.info({accounts: accountLevelGateways.map((m) => m.account_sid)},
|
||||
'multiple accounts have added this carrier with default routing -- cannot determine which to use');
|
||||
return {
|
||||
fromCarrier: true,
|
||||
error: 'Multiple accounts are attempting to default route this carrier'
|
||||
};
|
||||
}
|
||||
else if (accountLevelGateways.length === 1) {
|
||||
const [accounts] = await pp.query('SELECT * from accounts where account_sid = ?',
|
||||
accountLevelGateways[0].account_sid);
|
||||
return {
|
||||
fromCarrier: true,
|
||||
gateway: accountLevelGateways[0],
|
||||
service_provider_sid: accountLevelGateways[0].service_provider_sid,
|
||||
account_sid: accountLevelGateways[0].account_sid,
|
||||
application_sid: accountLevelGateways[0].application_sid,
|
||||
account: accounts[0]
|
||||
};
|
||||
}
|
||||
else {
|
||||
/* 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 || r[0].count > 1) 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],
|
||||
service_provider_sid: accounts[0].service_provider_sid,
|
||||
account_sid: accounts[0].account_sid,
|
||||
account: accounts[0]
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
else if (r.length > 1) {
|
||||
logger.info({r},
|
||||
'multiple accounts have added this carrier with default routing -- cannot determine which to use');
|
||||
return {
|
||||
fromCarrier: true,
|
||||
error: 'Multiple accounts are attempting to route the same phone number from the same carrier'
|
||||
};
|
||||
}
|
||||
|
||||
/* we have a route for this phone number and carrier combination */
|
||||
const gateway = matches.find((m) => m.voip_carrier_sid === r[0].voip_carrier_sid);
|
||||
const [accounts] = await pp.query(sqlAccountBySid, r[0].account_sid);
|
||||
assert(accounts.length);
|
||||
return {
|
||||
fromCarrier: true,
|
||||
gateway,
|
||||
service_provider_sid: accounts[0].service_provider_sid,
|
||||
account_sid: r[0].account_sid,
|
||||
application_sid: r[0].application_sid,
|
||||
account: accounts[0]
|
||||
};
|
||||
}
|
||||
}
|
||||
return failure;
|
||||
};
|
||||
|
||||
return {
|
||||
wasOriginatedFromCarrier,
|
||||
getApplicationForDidAndCarrier
|
||||
getApplicationForDidAndCarrier,
|
||||
getOutboundGatewayForRefer,
|
||||
getSPForAccount
|
||||
};
|
||||
};
|
||||
|
||||
+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}`);
|
||||
}
|
||||
|
||||
+169
-43
@@ -2,7 +2,7 @@ const debug = require('debug')('jambonz:sbc-inbound');
|
||||
const assert = require('assert');
|
||||
const Emitter = require('events');
|
||||
const parseUri = require('drachtio-srf').parseUri;
|
||||
const {makeCallCountKey} = require('./utils');
|
||||
const {nudgeCallCounts, roundTripTime} = require('./utils');
|
||||
const msProxyIps = process.env.MS_TEAMS_SIP_PROXY_IPS ?
|
||||
process.env.MS_TEAMS_SIP_PROXY_IPS.split(',').map((i) => i.trim()) :
|
||||
[];
|
||||
@@ -30,10 +30,10 @@ module.exports = function(srf, logger) {
|
||||
stats.histogram('app.hook.response_time', rtt, ['hook_type:auth', `status:${status}`]);
|
||||
})
|
||||
.on('error', async(err, req) => {
|
||||
const {account_sid} = req.locals;
|
||||
const {account_sid, account} = req.locals;
|
||||
const {writeAlerts, AlertType} = req.srf.locals;
|
||||
if (account_sid) {
|
||||
let opts = {account_sid};
|
||||
let opts = {account_sid, service_provider_sid: account.service_provider_sid};
|
||||
if (err.code === 'ECONNREFUSED') {
|
||||
opts = {...opts, alert_type: AlertType.WEBHOOK_CONNECTION_FAILURE, url: err.hook};
|
||||
}
|
||||
@@ -61,20 +61,30 @@ module.exports = function(srf, logger) {
|
||||
lookupAppByTeamsTenant,
|
||||
lookupAccountBySipRealm,
|
||||
lookupAccountBySid,
|
||||
lookupAccountCapacitiesBySid
|
||||
lookupAccountCapacitiesBySid,
|
||||
queryCallLimits
|
||||
} = srf.locals.dbHelpers;
|
||||
const {stats, writeCdrs} = srf.locals;
|
||||
const authenticator = require('@jambonz/http-authenticator')(lookupAuthHook, logger, {
|
||||
blacklistUnknownRealms: true,
|
||||
emitter: new AuthOutcomeReporter(stats)
|
||||
});
|
||||
const {wasOriginatedFromCarrier, getApplicationForDidAndCarrier} = require('./db-utils')(srf, logger);
|
||||
|
||||
|
||||
const initLocals = (req, res, next) => {
|
||||
req.locals = req.locals || {};
|
||||
req.locals.cdr = initCdr(req);
|
||||
const callId = req.get('Call-ID');
|
||||
req.locals = req.locals || {callId};
|
||||
|
||||
/* 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;
|
||||
}
|
||||
req.locals.cdr = initCdr(req);
|
||||
req.on('cancel', () => {
|
||||
logger.info({callId}, 'caller hungup before connecting to feature server');
|
||||
req.canceled = true;
|
||||
@@ -100,10 +110,45 @@ module.exports = function(srf, logger) {
|
||||
next();
|
||||
};
|
||||
|
||||
const handleSipRec = async(req, res, next) => {
|
||||
const {callId} = req.locals;
|
||||
if (Array.isArray(req.payload) && req.payload.length > 1) {
|
||||
const sdp = req.payload
|
||||
.find((p) => p.type === 'application/sdp')
|
||||
.content;
|
||||
if (!sdp) {
|
||||
logger.error({callId}, 'No SDP in multipart sdp');
|
||||
return res.send(503);
|
||||
}
|
||||
const xml = req.payload.find((p) => p.type !== 'application/sdp');
|
||||
const endPos = xml.content.indexOf('</recording>');
|
||||
xml.content = endPos !== -1 ?
|
||||
`${xml.content.substring(0, endPos + 12)}` :
|
||||
xml.content;
|
||||
logger.debug({callId, xml}, 'incoming call with SIPREC body');
|
||||
req.locals = {...req.locals, sdp, siprec: true, xml};
|
||||
}
|
||||
else req.locals = {...req.locals, sdp: req.body};
|
||||
next();
|
||||
};
|
||||
|
||||
const identifyAccount = async(req, res, next) => {
|
||||
try {
|
||||
|
||||
const {fromCarrier, gateway, account_sid, application_sid, account} = await wasOriginatedFromCarrier(req);
|
||||
const {siprec, callId} = req.locals;
|
||||
const {getSPForAccount, wasOriginatedFromCarrier, getApplicationForDidAndCarrier, stats} = req.srf.locals;
|
||||
const startAt = process.hrtime();
|
||||
const {
|
||||
fromCarrier,
|
||||
gateway,
|
||||
account_sid,
|
||||
application_sid,
|
||||
service_provider_sid,
|
||||
account,
|
||||
error
|
||||
} = await wasOriginatedFromCarrier(req);
|
||||
const rtt = roundTripTime(startAt);
|
||||
stats.histogram('app.mysql.response_time', rtt, [
|
||||
'query:wasOriginatedFromCarrier', 'app:sbc-inbound']);
|
||||
/**
|
||||
* calls come from 3 sources:
|
||||
* (1) A carrier
|
||||
@@ -111,18 +156,38 @@ module.exports = function(srf, logger) {
|
||||
* (3) A SIP user
|
||||
*/
|
||||
if (fromCarrier) {
|
||||
if (error) {
|
||||
return res.send(503, {
|
||||
headers: {
|
||||
'X-Reason': error
|
||||
}
|
||||
});
|
||||
}
|
||||
if (!gateway) {
|
||||
logger.info('identifyAccount: rejecting call from carrier because DID has not been provisioned');
|
||||
return res.send(404, 'Number Not Provisioned');
|
||||
}
|
||||
logger.debug({gateway}, 'identifyAccount: incoming call from gateway');
|
||||
|
||||
/* check for phone number level routing */
|
||||
const sid = application_sid || await getApplicationForDidAndCarrier(req, gateway.voip_carrier_sid);
|
||||
let sid;
|
||||
if (siprec) {
|
||||
if (!account.siprec_hook_sid) {
|
||||
logger.info({callId}, 'identifyAccount: rejecting call because SIPREC hook has not been provisioned');
|
||||
return res.send(404);
|
||||
}
|
||||
sid = account.siprec_hook_sid;
|
||||
}
|
||||
else {
|
||||
/* check for phone number level routing */
|
||||
sid = application_sid || await getApplicationForDidAndCarrier(req, gateway.voip_carrier_sid);
|
||||
}
|
||||
req.locals = {
|
||||
originator: 'trunk',
|
||||
carrier: gateway.name,
|
||||
gateway,
|
||||
voip_carrier_sid: gateway.voip_carrier_sid,
|
||||
application_sid: sid || gateway.application_sid,
|
||||
service_provider_sid,
|
||||
account_sid,
|
||||
account,
|
||||
...req.locals
|
||||
@@ -135,14 +200,16 @@ 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);
|
||||
}
|
||||
|
||||
const service_provider_sid = await getSPForAccount(app.account_sid);
|
||||
req.locals = {
|
||||
originator: 'teams',
|
||||
carrier: 'Microsoft Teams',
|
||||
msTeamsTenantFqdn: uri.host,
|
||||
account_sid: app.account_sid,
|
||||
service_provider_sid,
|
||||
...req.locals
|
||||
};
|
||||
}
|
||||
@@ -154,7 +221,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 */
|
||||
@@ -163,21 +231,26 @@ 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 = {
|
||||
service_provider_sid: account.service_provider_sid,
|
||||
account_sid: account.account_sid,
|
||||
account,
|
||||
application_sid: account.device_calling_application_sid,
|
||||
webhook_secret: account.webhook_secret,
|
||||
...req.locals
|
||||
};
|
||||
}
|
||||
assert(req.locals.service_provider_sid);
|
||||
assert(req.locals.account_sid);
|
||||
req.locals.cdr.account_sid = req.locals.account_sid;
|
||||
|
||||
if (!req.locals.account) {
|
||||
req.locals.account = await lookupAccountBySid(req.locals.account_sid);
|
||||
}
|
||||
req.locals.cdr.service_provider_sid = req.locals.account?.service_provider_sid;
|
||||
|
||||
if (!req.locals.account.is_active) {
|
||||
stats.increment('sbc.terminations', ['sipStatus:503']);
|
||||
@@ -189,7 +262,11 @@ module.exports = function(srf, logger) {
|
||||
delete req.locals.cdr;
|
||||
}
|
||||
|
||||
req.locals.logger = logger.child({callId: req.get('Call-ID'), account_sid: req.locals.account_sid});
|
||||
req.locals.logger = logger.child({
|
||||
callId: req.get('Call-ID'),
|
||||
service_provider_sid: req.locals.service_provider_sid,
|
||||
account_sid: req.locals.account_sid
|
||||
});
|
||||
|
||||
next();
|
||||
} catch (err) {
|
||||
@@ -200,30 +277,36 @@ module.exports = function(srf, logger) {
|
||||
};
|
||||
|
||||
const checkLimits = async(req, res, next) => {
|
||||
if (!process.env.JAMBONES_HOSTING) return next(); // skip
|
||||
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) return next(); // skip
|
||||
|
||||
const {incrKey, decrKey} = req.srf.locals.realtimeDbHelpers;
|
||||
const {logger, account_sid} = req.locals;
|
||||
const {writeAlerts, AlertType} = req.srf.locals;
|
||||
const {logger, account_sid, account, service_provider_sid, application_sid} = req.locals;
|
||||
const {writeCallCount, writeCallCountSP, writeCallCountApp, writeAlerts, AlertType} = req.srf.locals;
|
||||
assert(account_sid);
|
||||
const key = makeCallCountKey(account_sid);
|
||||
assert(service_provider_sid);
|
||||
|
||||
/* decrement count if INVITE is later rejected */
|
||||
res.once('end', ({status}) => {
|
||||
res.once('end', async({status}) => {
|
||||
if (status > 200) {
|
||||
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'));
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
/* increment the call count */
|
||||
const calls = await incrKey(key);
|
||||
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 ?
|
||||
@@ -232,25 +315,66 @@ module.exports = function(srf, logger) {
|
||||
logger.debug(`checkLimits: call count is now ${calls}, limit is ${minLimit}`);
|
||||
if (calls <= minLimit) return next();
|
||||
|
||||
const capacities = await lookupAccountCapacitiesBySid(account_sid);
|
||||
const limit = capacities.find((c) => c.category == 'voice_call_session');
|
||||
if (!limit) throw new Error('no account_capacities found');
|
||||
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'));
|
||||
return res.send(503, 'Maximum Calls In Progress');
|
||||
if (process.env.JAMBONES_HOSTING) {
|
||||
const accountCapacities = await lookupAccountCapacitiesBySid(account_sid);
|
||||
const accountLimit = accountCapacities.find((c) => c.category == 'voice_call_session');
|
||||
if (accountLimit) {
|
||||
/* check account limit */
|
||||
const limit_sessions = accountLimit.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.ACCOUNT_CALL_LIMIT,
|
||||
service_provider_sid: account.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);
|
||||
}
|
||||
}
|
||||
}
|
||||
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) {
|
||||
stats.increment('sbc.terminations', ['sipStatus:500']);
|
||||
logger.error({err}, 'error checking limits error for inbound call');
|
||||
res.send(500);
|
||||
req.srf.endSession(req);
|
||||
}
|
||||
};
|
||||
|
||||
@@ -263,11 +387,13 @@ 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);
|
||||
}
|
||||
};
|
||||
|
||||
return {
|
||||
initLocals,
|
||||
handleSipRec,
|
||||
challengeDeviceCalls,
|
||||
identifyAccount,
|
||||
checkLimits
|
||||
|
||||
+151
-8
@@ -12,19 +12,28 @@ const getAppserver = (srf) => {
|
||||
};
|
||||
|
||||
function makeRtpEngineOpts(req, srcIsUsingSrtp, dstIsUsingSrtp, teams = false) {
|
||||
const rtpCopy = JSON.parse(JSON.stringify(rtpCharacteristics));
|
||||
const srtpCopy = JSON.parse(JSON.stringify(srtpCharacteristics));
|
||||
const from = req.getParsedHeader('from');
|
||||
const srtpOpts = teams ? srtpCharacteristics['teams'] : srtpCharacteristics['default'];
|
||||
const dstOpts = dstIsUsingSrtp ? srtpOpts : rtpCharacteristics;
|
||||
const srctOpts = srcIsUsingSrtp ? srtpOpts : rtpCharacteristics;
|
||||
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']
|
||||
'replace': ['origin', 'session-connection'],
|
||||
'record call': process.env.JAMBONES_RECORD_ALL_CALLS ? 'yes' : 'no'
|
||||
};
|
||||
return {
|
||||
common,
|
||||
uas: {
|
||||
tag: from.params.tag,
|
||||
mediaOpts: srctOpts
|
||||
mediaOpts: srcOpts
|
||||
},
|
||||
uac: {
|
||||
tag: null,
|
||||
@@ -33,11 +42,16 @@ function makeRtpEngineOpts(req, srcIsUsingSrtp, dstIsUsingSrtp, teams = false) {
|
||||
};
|
||||
}
|
||||
|
||||
const SdpWantsSDES = (sdp) => {
|
||||
return /m=audio.*\s+RTP\/SAVP/.test(sdp);
|
||||
};
|
||||
const SdpWantsSrtp = (sdp) => {
|
||||
return /m=audio.*SAVP/.test(sdp);
|
||||
};
|
||||
|
||||
const makeCallCountKey = (sid) => `${sid}:incalls`;
|
||||
const makeAccountCallCountKey = (sid) => `incalls:account:${sid}`;
|
||||
const makeSPCallCountKey = (sid) => `incalls:sp:${sid}:`;
|
||||
const makeAppCallCountKey = (sid) => `incalls:app${sid}:`;
|
||||
|
||||
const normalizeDID = (tel) => {
|
||||
const regex = /^\+(\d+)$/;
|
||||
@@ -45,11 +59,140 @@ 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 doListen = (logger, app, port, resolve) => {
|
||||
return app.listen(port, () => {
|
||||
logger.info(`Health check server listening on http://localhost:${port}`);
|
||||
resolve(app);
|
||||
});
|
||||
};
|
||||
const handleErrors = (logger, app, resolve, reject, e) => {
|
||||
if (e.code === 'EADDRINUSE' &&
|
||||
process.env.HTTP_PORT_MAX &&
|
||||
e.port < process.env.HTTP_PORT_MAX) {
|
||||
|
||||
logger.info(`Health check server failed to bind port on ${e.port}, will try next port`);
|
||||
const server = doListen(logger, app, ++e.port, resolve);
|
||||
server.on('error', handleErrors.bind(null, logger, app, resolve, reject));
|
||||
return;
|
||||
}
|
||||
reject(e);
|
||||
};
|
||||
|
||||
|
||||
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, reject) => {
|
||||
const server = doListen(logger, app, port, resolve);
|
||||
server.on('error', handleErrors.bind(null, logger, app, resolve, reject));
|
||||
});
|
||||
};
|
||||
|
||||
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 = [];
|
||||
|
||||
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};
|
||||
};
|
||||
|
||||
const roundTripTime = (startAt) => {
|
||||
const diff = process.hrtime(startAt);
|
||||
const time = diff[0] * 1e3 + diff[1] * 1e-6;
|
||||
return time.toFixed(0);
|
||||
};
|
||||
|
||||
const parseConnectionIp = (sdp) => {
|
||||
const regex = /c=IN IP4 ([0-9.]+)/;
|
||||
const arr = regex.exec(sdp);
|
||||
return arr ? arr[1] : null;
|
||||
};
|
||||
|
||||
|
||||
module.exports = {
|
||||
isWSS,
|
||||
SdpWantsSrtp,
|
||||
SdpWantsSDES,
|
||||
getAppserver,
|
||||
makeRtpEngineOpts,
|
||||
makeCallCountKey,
|
||||
normalizeDID
|
||||
makeAccountCallCountKey,
|
||||
makeSPCallCountKey,
|
||||
makeAppCallCountKey,
|
||||
normalizeDID,
|
||||
equalsIgnoreOrder,
|
||||
systemHealth,
|
||||
createHealthCheckApp,
|
||||
nudgeCallCounts,
|
||||
roundTripTime,
|
||||
parseConnectionIp
|
||||
};
|
||||
|
||||
+24
@@ -0,0 +1,24 @@
|
||||
#!/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
|
||||
Generated
+5037
-2558
File diff suppressed because it is too large
Load Diff
+21
-21
@@ -1,9 +1,9 @@
|
||||
{
|
||||
"name": "sbc-inbound",
|
||||
"version": "0.3.6",
|
||||
"version": "v0.7.7",
|
||||
"main": "app.js",
|
||||
"engines": {
|
||||
"node": ">= 10.16.0"
|
||||
"node": ">= 12.0.0"
|
||||
},
|
||||
"keywords": [
|
||||
"sip",
|
||||
@@ -20,34 +20,34 @@
|
||||
},
|
||||
"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=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/ ",
|
||||
"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=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.12",
|
||||
"@jambonz/http-authenticator": "^0.2.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",
|
||||
"@jambonz/db-helpers": "^0.7.3",
|
||||
"@jambonz/http-authenticator": "^0.2.2",
|
||||
"@jambonz/http-health-check": "^0.0.1",
|
||||
"@jambonz/realtimedb-helpers": "^0.5.7",
|
||||
"@jambonz/rtpengine-utils": "^0.4.0",
|
||||
"@jambonz/siprec-client-utils": "^0.2.0",
|
||||
"@jambonz/stats-collector": "^0.1.6",
|
||||
"@jambonz/time-series": "^0.2.5",
|
||||
"aws-sdk": "^2.1261.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.4",
|
||||
"drachtio-fn-b2b-sugar": "0.0.12",
|
||||
"drachtio-srf": "^4.4.49",
|
||||
"pino": "^6.8.0",
|
||||
"rtpengine-client": "^0.2.0"
|
||||
"drachtio-srf": "^4.5.19",
|
||||
"express": "^4.18.1",
|
||||
"pino": "^7.11.0",
|
||||
"verify-aws-sns-signature": "^0.1.0",
|
||||
"xml2js": "^0.4.23"
|
||||
},
|
||||
"devDependencies": {
|
||||
"clear-module": "^4.1.1",
|
||||
"eslint": "^7.15.0",
|
||||
"eslint-plugin-promise": "^4.2.1",
|
||||
"eslint": "^7.32.0",
|
||||
"eslint-plugin-promise": "^4.3.1",
|
||||
"nyc": "^15.1.0",
|
||||
"tape": "^4.13.3"
|
||||
"tape": "^4.15.1"
|
||||
}
|
||||
}
|
||||
|
||||
+93
-18
@@ -4,6 +4,8 @@ SET FOREIGN_KEY_CHECKS=0;
|
||||
|
||||
DROP TABLE IF EXISTS account_static_ips;
|
||||
|
||||
DROP TABLE IF EXISTS account_limits;
|
||||
|
||||
DROP TABLE IF EXISTS account_products;
|
||||
|
||||
DROP TABLE IF EXISTS account_subscriptions;
|
||||
@@ -18,20 +20,28 @@ DROP TABLE IF EXISTS lcr_carrier_set_entry;
|
||||
|
||||
DROP TABLE IF EXISTS lcr_routes;
|
||||
|
||||
DROP TABLE IF EXISTS password_settings;
|
||||
|
||||
DROP TABLE IF EXISTS predefined_sip_gateways;
|
||||
|
||||
DROP TABLE IF EXISTS predefined_smpp_gateways;
|
||||
|
||||
DROP TABLE IF EXISTS predefined_carriers;
|
||||
|
||||
DROP TABLE IF EXISTS account_offers;
|
||||
|
||||
DROP TABLE IF EXISTS products;
|
||||
|
||||
DROP TABLE IF EXISTS schema_version;
|
||||
|
||||
DROP TABLE IF EXISTS api_keys;
|
||||
|
||||
DROP TABLE IF EXISTS sbc_addresses;
|
||||
|
||||
DROP TABLE IF EXISTS ms_teams_tenants;
|
||||
|
||||
DROP TABLE IF EXISTS service_provider_limits;
|
||||
|
||||
DROP TABLE IF EXISTS signup_history;
|
||||
|
||||
DROP TABLE IF EXISTS smpp_addresses;
|
||||
@@ -65,6 +75,15 @@ private_ipv4 VARBINARY(16) NOT NULL UNIQUE ,
|
||||
PRIMARY KEY (account_static_ip_sid)
|
||||
);
|
||||
|
||||
CREATE TABLE account_limits
|
||||
(
|
||||
account_limits_sid CHAR(36) NOT NULL UNIQUE ,
|
||||
account_sid CHAR(36) NOT NULL,
|
||||
category ENUM('api_rate','voice_call_session', 'device','voice_call_minutes','voice_call_session_license', 'voice_call_minutes_license') NOT NULL,
|
||||
quantity INTEGER NOT NULL,
|
||||
PRIMARY KEY (account_limits_sid)
|
||||
);
|
||||
|
||||
CREATE TABLE account_subscriptions
|
||||
(
|
||||
account_subscription_sid CHAR(36) NOT NULL UNIQUE ,
|
||||
@@ -119,6 +138,13 @@ priority INTEGER NOT NULL UNIQUE COMMENT 'lower priority routes are attempted f
|
||||
PRIMARY KEY (lcr_route_sid)
|
||||
) COMMENT='Least cost routing table';
|
||||
|
||||
CREATE TABLE password_settings
|
||||
(
|
||||
min_password_length INTEGER NOT NULL DEFAULT 8,
|
||||
require_digit BOOLEAN NOT NULL DEFAULT false,
|
||||
require_special_character BOOLEAN NOT NULL DEFAULT false
|
||||
);
|
||||
|
||||
CREATE TABLE predefined_carriers
|
||||
(
|
||||
predefined_carrier_sid CHAR(36) NOT NULL UNIQUE ,
|
||||
@@ -148,6 +174,20 @@ predefined_carrier_sid CHAR(36) NOT NULL,
|
||||
PRIMARY KEY (predefined_sip_gateway_sid)
|
||||
);
|
||||
|
||||
CREATE TABLE predefined_smpp_gateways
|
||||
(
|
||||
predefined_smpp_gateway_sid CHAR(36) NOT NULL UNIQUE ,
|
||||
ipv4 VARCHAR(128) NOT NULL COMMENT 'ip address or DNS name of the gateway. ',
|
||||
port INTEGER NOT NULL DEFAULT 2775 COMMENT 'smpp signaling port',
|
||||
inbound BOOLEAN NOT NULL COMMENT 'if true, whitelist this IP to allow inbound SMS from the gateway',
|
||||
outbound BOOLEAN NOT NULL COMMENT 'i',
|
||||
netmask INTEGER NOT NULL DEFAULT 32,
|
||||
is_primary BOOLEAN NOT NULL DEFAULT 1,
|
||||
use_tls BOOLEAN DEFAULT 0,
|
||||
predefined_carrier_sid CHAR(36) NOT NULL,
|
||||
PRIMARY KEY (predefined_smpp_gateway_sid)
|
||||
);
|
||||
|
||||
CREATE TABLE products
|
||||
(
|
||||
product_sid CHAR(36) NOT NULL UNIQUE ,
|
||||
@@ -174,6 +214,11 @@ stripe_product_id VARCHAR(56) NOT NULL,
|
||||
PRIMARY KEY (account_offer_sid)
|
||||
);
|
||||
|
||||
CREATE TABLE schema_version
|
||||
(
|
||||
version VARCHAR(16)
|
||||
);
|
||||
|
||||
CREATE TABLE api_keys
|
||||
(
|
||||
api_key_sid CHAR(36) NOT NULL UNIQUE ,
|
||||
@@ -205,6 +250,15 @@ tenant_fqdn VARCHAR(255) NOT NULL UNIQUE ,
|
||||
PRIMARY KEY (ms_teams_tenant_sid)
|
||||
) COMMENT='A Microsoft Teams customer tenant';
|
||||
|
||||
CREATE TABLE service_provider_limits
|
||||
(
|
||||
service_provider_limits_sid CHAR(36) NOT NULL UNIQUE ,
|
||||
service_provider_sid CHAR(36) NOT NULL,
|
||||
category ENUM('api_rate','voice_call_session', 'device','voice_call_minutes','voice_call_session_license', 'voice_call_minutes_license') NOT NULL,
|
||||
quantity INTEGER NOT NULL,
|
||||
PRIMARY KEY (service_provider_limits_sid)
|
||||
);
|
||||
|
||||
CREATE TABLE signup_history
|
||||
(
|
||||
email VARCHAR(255) NOT NULL,
|
||||
@@ -260,6 +314,7 @@ email_activation_code VARCHAR(16),
|
||||
email_validated BOOLEAN NOT NULL DEFAULT false,
|
||||
phone_validated BOOLEAN NOT NULL DEFAULT false,
|
||||
email_content_opt_out BOOLEAN NOT NULL DEFAULT false,
|
||||
is_active BOOLEAN NOT NULL DEFAULT true,
|
||||
PRIMARY KEY (user_sid)
|
||||
);
|
||||
|
||||
@@ -287,6 +342,9 @@ smpp_password VARCHAR(64),
|
||||
smpp_enquire_link_interval INTEGER DEFAULT 0,
|
||||
smpp_inbound_system_id VARCHAR(255),
|
||||
smpp_inbound_password VARCHAR(64),
|
||||
register_from_user VARCHAR(128),
|
||||
register_from_domain VARCHAR(255),
|
||||
register_public_ip_in_contact BOOLEAN NOT NULL DEFAULT false,
|
||||
PRIMARY KEY (voip_carrier_sid)
|
||||
) COMMENT='A Carrier or customer PBX that can send or receive calls';
|
||||
|
||||
@@ -395,6 +453,11 @@ disable_cdrs BOOLEAN NOT NULL DEFAULT 0,
|
||||
trial_end_date DATETIME,
|
||||
deactivated_reason VARCHAR(255),
|
||||
device_to_call_ratio INTEGER NOT NULL DEFAULT 5,
|
||||
subspace_client_id VARCHAR(255),
|
||||
subspace_client_secret VARCHAR(255),
|
||||
subspace_sip_teleport_id VARCHAR(255),
|
||||
subspace_sip_teleport_destinations VARCHAR(255),
|
||||
siprec_hook_sid CHAR(36),
|
||||
PRIMARY KEY (account_sid)
|
||||
) COMMENT='An enterprise that uses the platform for comm services';
|
||||
|
||||
@@ -402,24 +465,31 @@ CREATE INDEX account_static_ip_sid_idx ON account_static_ips (account_static_ip_
|
||||
CREATE INDEX account_sid_idx ON account_static_ips (account_sid);
|
||||
ALTER TABLE account_static_ips ADD FOREIGN KEY account_sid_idxfk (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
CREATE INDEX account_sid_idx ON account_limits (account_sid);
|
||||
ALTER TABLE account_limits ADD FOREIGN KEY account_sid_idxfk_1 (account_sid) REFERENCES accounts (account_sid) ON DELETE CASCADE;
|
||||
|
||||
CREATE INDEX account_subscription_sid_idx ON account_subscriptions (account_subscription_sid);
|
||||
CREATE INDEX account_sid_idx ON account_subscriptions (account_sid);
|
||||
ALTER TABLE account_subscriptions ADD FOREIGN KEY account_sid_idxfk_1 (account_sid) REFERENCES accounts (account_sid);
|
||||
ALTER TABLE account_subscriptions ADD FOREIGN KEY account_sid_idxfk_2 (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
CREATE INDEX invite_code_idx ON beta_invite_codes (invite_code);
|
||||
CREATE INDEX call_route_sid_idx ON call_routes (call_route_sid);
|
||||
ALTER TABLE call_routes ADD FOREIGN KEY account_sid_idxfk_2 (account_sid) REFERENCES accounts (account_sid);
|
||||
ALTER TABLE call_routes ADD FOREIGN KEY account_sid_idxfk_3 (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
ALTER TABLE call_routes ADD FOREIGN KEY application_sid_idxfk (application_sid) REFERENCES applications (application_sid);
|
||||
|
||||
CREATE INDEX dns_record_sid_idx ON dns_records (dns_record_sid);
|
||||
ALTER TABLE dns_records ADD FOREIGN KEY account_sid_idxfk_3 (account_sid) REFERENCES accounts (account_sid);
|
||||
ALTER TABLE dns_records ADD FOREIGN KEY account_sid_idxfk_4 (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
CREATE INDEX predefined_carrier_sid_idx ON predefined_carriers (predefined_carrier_sid);
|
||||
CREATE INDEX predefined_sip_gateway_sid_idx ON predefined_sip_gateways (predefined_sip_gateway_sid);
|
||||
CREATE INDEX predefined_carrier_sid_idx ON predefined_sip_gateways (predefined_carrier_sid);
|
||||
ALTER TABLE predefined_sip_gateways ADD FOREIGN KEY predefined_carrier_sid_idxfk (predefined_carrier_sid) REFERENCES predefined_carriers (predefined_carrier_sid);
|
||||
|
||||
CREATE INDEX predefined_smpp_gateway_sid_idx ON predefined_smpp_gateways (predefined_smpp_gateway_sid);
|
||||
CREATE INDEX predefined_carrier_sid_idx ON predefined_smpp_gateways (predefined_carrier_sid);
|
||||
ALTER TABLE predefined_smpp_gateways ADD FOREIGN KEY predefined_carrier_sid_idxfk_1 (predefined_carrier_sid) REFERENCES predefined_carriers (predefined_carrier_sid);
|
||||
|
||||
CREATE INDEX product_sid_idx ON products (product_sid);
|
||||
CREATE INDEX account_product_sid_idx ON account_products (account_product_sid);
|
||||
CREATE INDEX account_subscription_sid_idx ON account_products (account_subscription_sid);
|
||||
@@ -429,14 +499,14 @@ ALTER TABLE account_products ADD FOREIGN KEY product_sid_idxfk (product_sid) REF
|
||||
|
||||
CREATE INDEX account_offer_sid_idx ON account_offers (account_offer_sid);
|
||||
CREATE INDEX account_sid_idx ON account_offers (account_sid);
|
||||
ALTER TABLE account_offers ADD FOREIGN KEY account_sid_idxfk_4 (account_sid) REFERENCES accounts (account_sid);
|
||||
ALTER TABLE account_offers ADD FOREIGN KEY account_sid_idxfk_5 (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
CREATE INDEX product_sid_idx ON account_offers (product_sid);
|
||||
ALTER TABLE account_offers ADD FOREIGN KEY product_sid_idxfk_1 (product_sid) REFERENCES products (product_sid);
|
||||
|
||||
CREATE INDEX api_key_sid_idx ON api_keys (api_key_sid);
|
||||
CREATE INDEX account_sid_idx ON api_keys (account_sid);
|
||||
ALTER TABLE api_keys ADD FOREIGN KEY account_sid_idxfk_5 (account_sid) REFERENCES accounts (account_sid);
|
||||
ALTER TABLE api_keys ADD FOREIGN KEY account_sid_idxfk_6 (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
CREATE INDEX service_provider_sid_idx ON api_keys (service_provider_sid);
|
||||
ALTER TABLE api_keys ADD FOREIGN KEY service_provider_sid_idxfk (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
@@ -450,41 +520,44 @@ ALTER TABLE sbc_addresses ADD FOREIGN KEY service_provider_sid_idxfk_1 (service_
|
||||
CREATE INDEX ms_teams_tenant_sid_idx ON ms_teams_tenants (ms_teams_tenant_sid);
|
||||
ALTER TABLE ms_teams_tenants ADD FOREIGN KEY service_provider_sid_idxfk_2 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
|
||||
ALTER TABLE ms_teams_tenants ADD FOREIGN KEY account_sid_idxfk_6 (account_sid) REFERENCES accounts (account_sid);
|
||||
ALTER TABLE ms_teams_tenants ADD FOREIGN KEY account_sid_idxfk_7 (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
ALTER TABLE ms_teams_tenants ADD FOREIGN KEY application_sid_idxfk_1 (application_sid) REFERENCES applications (application_sid);
|
||||
|
||||
CREATE INDEX tenant_fqdn_idx ON ms_teams_tenants (tenant_fqdn);
|
||||
CREATE INDEX service_provider_sid_idx ON service_provider_limits (service_provider_sid);
|
||||
ALTER TABLE service_provider_limits ADD FOREIGN KEY service_provider_sid_idxfk_3 (service_provider_sid) REFERENCES service_providers (service_provider_sid) ON DELETE CASCADE;
|
||||
|
||||
CREATE INDEX email_idx ON signup_history (email);
|
||||
CREATE INDEX smpp_address_sid_idx ON smpp_addresses (smpp_address_sid);
|
||||
CREATE INDEX service_provider_sid_idx ON smpp_addresses (service_provider_sid);
|
||||
ALTER TABLE smpp_addresses ADD FOREIGN KEY service_provider_sid_idxfk_3 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
ALTER TABLE smpp_addresses ADD FOREIGN KEY service_provider_sid_idxfk_4 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
|
||||
CREATE UNIQUE INDEX speech_credentials_idx_1 ON speech_credentials (vendor,account_sid);
|
||||
|
||||
CREATE INDEX speech_credential_sid_idx ON speech_credentials (speech_credential_sid);
|
||||
CREATE INDEX service_provider_sid_idx ON speech_credentials (service_provider_sid);
|
||||
ALTER TABLE speech_credentials ADD FOREIGN KEY service_provider_sid_idxfk_4 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
ALTER TABLE speech_credentials ADD FOREIGN KEY service_provider_sid_idxfk_5 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
|
||||
CREATE INDEX account_sid_idx ON speech_credentials (account_sid);
|
||||
ALTER TABLE speech_credentials ADD FOREIGN KEY account_sid_idxfk_7 (account_sid) REFERENCES accounts (account_sid);
|
||||
ALTER TABLE speech_credentials ADD FOREIGN KEY account_sid_idxfk_8 (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
CREATE INDEX user_sid_idx ON users (user_sid);
|
||||
CREATE INDEX email_idx ON users (email);
|
||||
CREATE INDEX phone_idx ON users (phone);
|
||||
CREATE INDEX account_sid_idx ON users (account_sid);
|
||||
ALTER TABLE users ADD FOREIGN KEY account_sid_idxfk_8 (account_sid) REFERENCES accounts (account_sid);
|
||||
ALTER TABLE users ADD FOREIGN KEY account_sid_idxfk_9 (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
CREATE INDEX service_provider_sid_idx ON users (service_provider_sid);
|
||||
ALTER TABLE users ADD FOREIGN KEY service_provider_sid_idxfk_5 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
ALTER TABLE users ADD FOREIGN KEY service_provider_sid_idxfk_6 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
|
||||
CREATE INDEX email_activation_code_idx ON users (email_activation_code);
|
||||
CREATE INDEX voip_carrier_sid_idx ON voip_carriers (voip_carrier_sid);
|
||||
CREATE INDEX account_sid_idx ON voip_carriers (account_sid);
|
||||
ALTER TABLE voip_carriers ADD FOREIGN KEY account_sid_idxfk_9 (account_sid) REFERENCES accounts (account_sid);
|
||||
ALTER TABLE voip_carriers ADD FOREIGN KEY account_sid_idxfk_10 (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
CREATE INDEX service_provider_sid_idx ON voip_carriers (service_provider_sid);
|
||||
ALTER TABLE voip_carriers ADD FOREIGN KEY service_provider_sid_idxfk_6 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
ALTER TABLE voip_carriers ADD FOREIGN KEY service_provider_sid_idxfk_7 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
|
||||
ALTER TABLE voip_carriers ADD FOREIGN KEY application_sid_idxfk_2 (application_sid) REFERENCES applications (application_sid);
|
||||
|
||||
@@ -497,12 +570,12 @@ CREATE INDEX number_idx ON phone_numbers (number);
|
||||
CREATE INDEX voip_carrier_sid_idx ON phone_numbers (voip_carrier_sid);
|
||||
ALTER TABLE phone_numbers ADD FOREIGN KEY voip_carrier_sid_idxfk_1 (voip_carrier_sid) REFERENCES voip_carriers (voip_carrier_sid);
|
||||
|
||||
ALTER TABLE phone_numbers ADD FOREIGN KEY account_sid_idxfk_10 (account_sid) REFERENCES accounts (account_sid);
|
||||
ALTER TABLE phone_numbers ADD FOREIGN KEY account_sid_idxfk_11 (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
ALTER TABLE phone_numbers ADD FOREIGN KEY application_sid_idxfk_3 (application_sid) REFERENCES applications (application_sid);
|
||||
|
||||
CREATE INDEX service_provider_sid_idx ON phone_numbers (service_provider_sid);
|
||||
ALTER TABLE phone_numbers ADD FOREIGN KEY service_provider_sid_idxfk_7 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
ALTER TABLE phone_numbers ADD FOREIGN KEY service_provider_sid_idxfk_8 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
|
||||
CREATE INDEX sip_gateway_idx_hostport ON sip_gateways (ipv4,port);
|
||||
|
||||
@@ -518,10 +591,10 @@ CREATE UNIQUE INDEX applications_idx_name ON applications (account_sid,name);
|
||||
|
||||
CREATE INDEX application_sid_idx ON applications (application_sid);
|
||||
CREATE INDEX service_provider_sid_idx ON applications (service_provider_sid);
|
||||
ALTER TABLE applications ADD FOREIGN KEY service_provider_sid_idxfk_8 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
ALTER TABLE applications ADD FOREIGN KEY service_provider_sid_idxfk_9 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
|
||||
CREATE INDEX account_sid_idx ON applications (account_sid);
|
||||
ALTER TABLE applications ADD FOREIGN KEY account_sid_idxfk_11 (account_sid) REFERENCES accounts (account_sid);
|
||||
ALTER TABLE applications ADD FOREIGN KEY account_sid_idxfk_12 (account_sid) REFERENCES accounts (account_sid);
|
||||
|
||||
ALTER TABLE applications ADD FOREIGN KEY call_hook_sid_idxfk (call_hook_sid) REFERENCES webhooks (webhook_sid);
|
||||
|
||||
@@ -537,7 +610,7 @@ ALTER TABLE service_providers ADD FOREIGN KEY registration_hook_sid_idxfk (regis
|
||||
CREATE INDEX account_sid_idx ON accounts (account_sid);
|
||||
CREATE INDEX sip_realm_idx ON accounts (sip_realm);
|
||||
CREATE INDEX service_provider_sid_idx ON accounts (service_provider_sid);
|
||||
ALTER TABLE accounts ADD FOREIGN KEY service_provider_sid_idxfk_9 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
ALTER TABLE accounts ADD FOREIGN KEY service_provider_sid_idxfk_10 (service_provider_sid) REFERENCES service_providers (service_provider_sid);
|
||||
|
||||
ALTER TABLE accounts ADD FOREIGN KEY registration_hook_sid_idxfk_1 (registration_hook_sid) REFERENCES webhooks (webhook_sid);
|
||||
|
||||
@@ -545,4 +618,6 @@ ALTER TABLE accounts ADD FOREIGN KEY queue_event_hook_sid_idxfk (queue_event_hoo
|
||||
|
||||
ALTER TABLE accounts ADD FOREIGN KEY device_calling_application_sid_idxfk (device_calling_application_sid) REFERENCES applications (application_sid);
|
||||
|
||||
ALTER TABLE accounts ADD FOREIGN KEY siprec_hook_sid_idxfk (siprec_hook_sid) REFERENCES applications (application_sid);
|
||||
|
||||
SET FOREIGN_KEY_CHECKS=1;
|
||||
|
||||
@@ -9,6 +9,8 @@ insert into webhooks(webhook_sid, url, username, password) values('90dda62e-0ea2
|
||||
|
||||
insert into service_providers (service_provider_sid, name, root_domain, registration_hook_sid)
|
||||
values ('3f35518f-5a0d-4c2e-90a5-2407bb3b36f0', 'SP A', 'jambonz.org', '90dda62e-0ea2-47d1-8164-5bd49003476c');
|
||||
insert into service_provider_limits (service_provider_limits_sid, service_provider_sid, category, quantity)
|
||||
values ('a79d3ade-e0da-4461-80f3-7c73f01e18b4', '3f35518f-5a0d-4c2e-90a5-2407bb3b36f0', 'voice_call_session', 1);
|
||||
|
||||
insert into accounts(account_sid, service_provider_sid, name, sip_realm, registration_hook_sid, webhook_secret)
|
||||
values ('ed649e33-e771-403a-8c99-1780eabbc803', '3f35518f-5a0d-4c2e-90a5-2407bb3b36f0', 'test account', 'jambonz.org', '90dda62e-0ea2-47d1-8164-5bd49003476c', 'foobar');
|
||||
@@ -38,6 +40,7 @@ insert into account_subscriptions(account_subscription_sid, account_sid, pending
|
||||
values ('73bbcc5d-512f-4cea-8535-9a6e3d2bd19d','d7cc37cb-d152-49ef-a51b-485f6e917089',0);
|
||||
insert into account_products(account_product_sid, account_subscription_sid, product_sid,quantity)
|
||||
values ('92f137f7-4bc3-4157-b096-6817e54b1874', '73bbcc5d-512f-4cea-8535-9a6e3d2bd19d', 'c4403cdb-8e75-4b27-9726-7d8315e3216d', 0);
|
||||
insert into account_limits(account_limits_sid, account_sid, category, quantity) values('a1b2c3d4-e5f6-7a8b-9c0d-1e2f3a4b5c6d', 'd7cc37cb-d152-49ef-a51b-485f6e917089', 'voice_call_session', 0);
|
||||
|
||||
insert into voip_carriers (voip_carrier_sid, name, account_sid) values ('9b1abdc7-0220-4964-bc66-32b5c70cd9ab', 'westco', 'd7cc37cb-d152-49ef-a51b-485f6e917089');
|
||||
insert into sip_gateways (sip_gateway_sid, voip_carrier_sid, ipv4, inbound, outbound)
|
||||
@@ -57,4 +60,32 @@ insert into phone_numbers (phone_number_sid, number, voip_carrier_sid, account_s
|
||||
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 ('f7ad205d-b92f-4363-8160-f8b5216b40d3', '15083871234', '287c1452-620d-4195-9f19-c9814ef90d78', 'd7cc37cb-d152-49ef-a51b-485f6e917089');
|
||||
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');
|
||||
|
||||
-- two accounts that both have the same carrier with default routing (ambiguity test)
|
||||
insert into accounts (account_sid, name, service_provider_sid, webhook_secret, sip_realm)
|
||||
values ('239d7d49-b3e4-4fdb-9d66-661149f717e8', 'Account B1', '3f35518f-5a0d-4c2e-90a5-2407bb3b36f0', 'foobar', 'echo2.sip.jambonz.org');
|
||||
insert into accounts (account_sid, name, service_provider_sid, webhook_secret, sip_realm)
|
||||
values ('909d7d49-b3e4-4fdb-9d66-661149f717e8', 'Account B2', '3f35518f-5a0d-4c2e-90a5-2407bb3b36f0', 'foobar', 'foxtrot.sip.jambonz.org');
|
||||
|
||||
insert into applications (application_sid, name, account_sid, call_hook_sid, call_status_hook_sid)
|
||||
values ('8843e39f-4346-4218-8434-a53130e8be49', 'test', '239d7d49-b3e4-4fdb-9d66-661149f717e8', '90dda62e-0ea2-47d1-8164-5bd49003476c', '4d7ce0aa-5ead-4e61-9a6b-3daa732218b1');
|
||||
insert into applications (application_sid, name, account_sid, call_hook_sid, call_status_hook_sid)
|
||||
values ('7743e39f-4346-4218-8434-a53130e8be49', 'test', '909d7d49-b3e4-4fdb-9d66-661149f717e8', '90dda62e-0ea2-47d1-8164-5bd49003476c', '4d7ce0aa-5ead-4e61-9a6b-3daa732218b1');
|
||||
|
||||
insert into voip_carriers (voip_carrier_sid, name, service_provider_sid, account_sid, application_sid)
|
||||
values ('731abdc7-0220-4964-bc66-32b5c70cd9ab', 'twilio-1', '3f35518f-5a0d-4c2e-90a5-2407bb3b36f0', '239d7d49-b3e4-4fdb-9d66-661149f717e8', '8843e39f-4346-4218-8434-a53130e8be49');
|
||||
insert into voip_carriers (voip_carrier_sid, name, service_provider_sid, account_sid, application_sid)
|
||||
values ('987abdc7-0220-4964-bc66-32b5c70cd9ab', 'twilio-2', '3f35518f-5a0d-4c2e-90a5-2407bb3b36f0', '909d7d49-b3e4-4fdb-9d66-661149f717e8', '7743e39f-4346-4218-8434-a53130e8be49');
|
||||
|
||||
insert into sip_gateways (sip_gateway_sid, voip_carrier_sid, ipv4, inbound, outbound)
|
||||
values ('664a5339-c62c-4075-9e19-f4de70a96597', '731abdc7-0220-4964-bc66-32b5c70cd9ab', '172.38.0.40', true, false);
|
||||
insert into sip_gateways (sip_gateway_sid, voip_carrier_sid, ipv4, inbound, outbound)
|
||||
values ('554a5339-c62c-4075-9e19-f4de70a96597', '987abdc7-0220-4964-bc66-32b5c70cd9ab', '172.38.0.40', true, false);
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -0,0 +1,118 @@
|
||||
<?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>
|
||||
@@ -0,0 +1,70 @@
|
||||
<?xml version="1.0" encoding="ISO-8859-1" ?>
|
||||
<!DOCTYPE scenario SYSTEM "sipp.dtd">
|
||||
|
||||
<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:+15083871234@172.38.0.10 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:15083871234@172.38.0.10>
|
||||
Call-ID: [call_id]
|
||||
CSeq: 1 INVITE
|
||||
Contact: sip:sipp@[local_ip]:[local_port]
|
||||
Max-Forwards: 70
|
||||
Subject: uac-pcap-carrier-fail-ambiguous
|
||||
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>
|
||||
|
||||
|
||||
<!-- 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="503" 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:15083871234@172.38.0.10 SIP/2.0
|
||||
[last_Via]
|
||||
From: sipp <sip:sipp@[local_ip]:[local_port]>;tag=[pid]SIPpTag09[call_number]
|
||||
To: <sip:15083871234@172.38.0.10>[peer_tag_param]
|
||||
Call-ID: [call_id]
|
||||
CSeq: 1 ACK
|
||||
Max-Forwards: 70
|
||||
Subject: uac-pcap-carrier-fail-ambiguous
|
||||
Content-Length: 0
|
||||
|
||||
]]>
|
||||
</send>
|
||||
|
||||
|
||||
<!-- 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>
|
||||
+18
-1
@@ -124,7 +124,24 @@
|
||||
]]>
|
||||
</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>
|
||||
|
||||
|
||||
+22
-9
@@ -1,8 +1,7 @@
|
||||
const test = require('tape');
|
||||
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};
|
||||
const { sippUac } = require('./sipp')('test_sbc-inbound');
|
||||
const bent = require('bent');
|
||||
const getJSON = bent('json');
|
||||
|
||||
process.on('unhandledRejection', (reason, p) => {
|
||||
console.log('Unhandled Rejection at: Promise', p, 'reason:', reason);
|
||||
@@ -25,15 +24,24 @@ function waitFor(ms) {
|
||||
test('incoming call tests', async(t) => {
|
||||
const {srf} = require('../app');
|
||||
const { queryCdrs } = srf.locals;
|
||||
let res;
|
||||
|
||||
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');
|
||||
|
||||
@@ -50,12 +58,17 @@ test('incoming call tests', async(t) => {
|
||||
t.pass('handles in-dialog requests');
|
||||
|
||||
await sippUac('uac-pcap-carrier-max-call-limit.xml', '172.38.0.20');
|
||||
t.pass('rejects incoming call with 503 when max calls reached')
|
||||
t.pass('rejects incoming call with 503 when max calls per account reached');
|
||||
|
||||
/* switch off this env for remaining tests (JAMBONES_HOSTING is for Saas sts) */
|
||||
delete process.env.JAMBONES_HOSTING;
|
||||
await sippUac('uac-pcap-carrier-fail-ambiguous.xml', '172.38.0.40');
|
||||
t.pass('rejects incoming call with 503 when multiple accounts have same carrier witrh default routing')
|
||||
|
||||
await waitFor(10);
|
||||
await waitFor(12);
|
||||
const res = await queryCdrs({account_sid: 'ed649e33-e771-403a-8c99-1780eabbc803'});
|
||||
console.log(`cdrs: ${JSON.stringify(res)}`);
|
||||
t.ok(6 === res.total, 'successfully wrote 6 cdrs for calls');
|
||||
//console.log(`cdrs: ${JSON.stringify(res)}`);
|
||||
t.ok(7 === res.total, 'successfully wrote 7 cdrs for calls');
|
||||
|
||||
srf.disconnect();
|
||||
t.end();
|
||||
|
||||
Reference in New Issue
Block a user