mirror of
https://github.com/jambonz/sbc-inbound.git
synced 2026-10-04 02:04:22 +00:00
Compare commits
32
Commits
v0.6.4
..
v0.7.2-rc5
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
75d8381ebb | ||
|
|
7750fdc3f2 | ||
|
|
fcf90916e1 | ||
|
|
4a0aa79456 | ||
|
|
1cdaae7773 | ||
|
|
82ac86389d | ||
|
|
e9184ad208 | ||
|
|
804bb890b5 | ||
|
|
a892a87eb5 | ||
|
|
56205dc852 | ||
|
|
927b7be637 | ||
|
|
63b482e562 | ||
|
|
a7d047a7e8 | ||
|
|
ea4f5ea0a8 | ||
|
|
13d8ee8c3c | ||
|
|
c717470c5e | ||
|
|
f16598e144 | ||
|
|
8584020d4c | ||
|
|
0b1b49bd59 | ||
|
|
e22b2ae63f | ||
|
|
2a5200673d | ||
|
|
13ca308df7 | ||
|
|
981788dc58 | ||
|
|
ca4599da22 | ||
|
|
04c4ad5975 | ||
|
|
5267b6f9a4 | ||
|
|
2fa1f08244 | ||
|
|
03cb1cceb5 | ||
|
|
9b78a818b5 | ||
|
|
e6a90a32ad | ||
|
|
fb21fd2333 | ||
|
|
819fecf455 |
@@ -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: 14.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
|
||||
|
||||
+2
-8
@@ -1,16 +1,10 @@
|
||||
FROM node:alpine as builder
|
||||
RUN apk update && apk add --no-cache python make g++
|
||||
FROM node:17-slim
|
||||
WORKDIR /opt/app/
|
||||
COPY package.json ./
|
||||
RUN npm install
|
||||
RUN npm prune
|
||||
|
||||
FROM node:alpine as app
|
||||
WORKDIR /opt/app
|
||||
COPY . /opt/app
|
||||
COPY --from=builder /opt/app/node_modules ./node_modules
|
||||
|
||||
ARG NODE_ENV
|
||||
ENV NODE_ENV $NODE_ENV
|
||||
|
||||
CMD [ "npm", "start" ]
|
||||
CMD [ "npm", "start" ]
|
||||
@@ -1,6 +1,6 @@
|
||||
MIT License
|
||||
|
||||
Copyright (c) 2019 jambonz
|
||||
Copyright (c) 2021 Drachtio Communications Services, LLC
|
||||
|
||||
Permission is hereby granted, free of charge, to any person obtaining a copy
|
||||
of this software and associated documentation files (the "Software"), to deal
|
||||
|
||||
@@ -6,11 +6,9 @@ assert.ok(process.env.JAMBONES_MYSQL_HOST &&
|
||||
assert.ok(process.env.DRACHTIO_PORT || process.env.DRACHTIO_HOST, 'missing DRACHTIO_PORT env var');
|
||||
assert.ok(process.env.DRACHTIO_SECRET, 'missing DRACHTIO_SECRET env var');
|
||||
assert.ok(process.env.JAMBONES_TIME_SERIES_HOST, 'missing JAMBONES_TIME_SERIES_HOST env var');
|
||||
assert.ok(process.env.JAMBONES_NETWORK_CIDR, 'missing JAMBONES_NETWORK_CIDR env var');
|
||||
assert.ok(process.env.JAMBONES_NETWORK_CIDR || process.env.K8S, 'missing JAMBONES_NETWORK_CIDR env var');
|
||||
const Srf = require('drachtio-srf');
|
||||
const srf = new Srf('sbc-inbound');
|
||||
const CIDRMatcher = require('cidr-matcher');
|
||||
const matcher = new CIDRMatcher([process.env.JAMBONES_NETWORK_CIDR]);
|
||||
const opts = Object.assign({
|
||||
timestamp: () => {return `, "time": "${new Date().toISOString()}"`;}
|
||||
}, {level: process.env.JAMBONES_LOGLEVEL || 'info'});
|
||||
@@ -22,13 +20,17 @@ const {
|
||||
AlertType
|
||||
} = require('@jambonz/time-series')(logger, {
|
||||
host: process.env.JAMBONES_TIME_SERIES_HOST,
|
||||
port: process.env.JAMBONES_TIME_SERIES_PORT || 8086,
|
||||
commitSize: 50,
|
||||
commitInterval: 'test' === process.env.NODE_ENV ? 7 : 20
|
||||
});
|
||||
const StatsCollector = require('@jambonz/stats-collector');
|
||||
const stats = new StatsCollector(logger);
|
||||
const {equalsIgnoreOrder} = require('./lib/utils');
|
||||
const {LifeCycleEvents} = require('./lib/constants');
|
||||
const setNameRtp = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-rtp`;
|
||||
const rtpServers = [];
|
||||
const setName = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-sip`;
|
||||
|
||||
const {
|
||||
pool,
|
||||
@@ -47,12 +49,16 @@ const {
|
||||
database: process.env.JAMBONES_MYSQL_DATABASE,
|
||||
connectionLimit: process.env.JAMBONES_MYSQL_CONNECTION_LIMIT || 10
|
||||
}, logger);
|
||||
const {createSet, retrieveSet, incrKey, decrKey} = require('@jambonz/realtimedb-helpers')({
|
||||
const {createSet, retrieveSet, addToSet, removeFromSet, incrKey, decrKey} = require('@jambonz/realtimedb-helpers')({
|
||||
host: process.env.JAMBONES_REDIS_HOST || 'localhost',
|
||||
port: process.env.JAMBONES_REDIS_PORT || 6379
|
||||
}, logger);
|
||||
|
||||
const {getRtpEngine, setRtpEngines} = require('@jambonz/rtpengine-utils')([], logger, {emitter: stats});
|
||||
const {getRtpEngine, setRtpEngines} = require('@jambonz/rtpengine-utils')([], logger, {
|
||||
emitter: stats,
|
||||
dtmfListenPort: process.env.DTMF_LISTEN_PORT || 22224,
|
||||
protocol: 'udp'
|
||||
});
|
||||
srf.locals = {...srf.locals,
|
||||
stats,
|
||||
queryCdrs,
|
||||
@@ -78,7 +84,18 @@ srf.locals = {...srf.locals,
|
||||
retrieveSet
|
||||
}
|
||||
};
|
||||
srf.locals.getFeatureServer = require('./lib/fs-tracking')(srf, logger);
|
||||
const {
|
||||
wasOriginatedFromCarrier,
|
||||
getApplicationForDidAndCarrier,
|
||||
getOutboundGatewayForRefer
|
||||
} = require('./lib/db-utils')(srf, logger);
|
||||
srf.locals = {
|
||||
...srf.locals,
|
||||
wasOriginatedFromCarrier,
|
||||
getApplicationForDidAndCarrier,
|
||||
getOutboundGatewayForRefer,
|
||||
getFeatureServer: require('./lib/fs-tracking')(srf, logger)
|
||||
};
|
||||
const activeCallIds = srf.locals.activeCallIds;
|
||||
|
||||
const {
|
||||
@@ -89,25 +106,39 @@ const {
|
||||
} = require('./lib/middleware')(srf, logger);
|
||||
const CallSession = require('./lib/call-session');
|
||||
|
||||
if (process.env.DRACHTIO_HOST) {
|
||||
if (process.env.DRACHTIO_HOST && !process.env.K8S) {
|
||||
const CIDRMatcher = require('cidr-matcher');
|
||||
const cidrs = process.env.JAMBONES_NETWORK_CIDR
|
||||
.split(',')
|
||||
.map((s) => s.trim());
|
||||
const matcher = new CIDRMatcher(cidrs);
|
||||
|
||||
srf.connect({host: process.env.DRACHTIO_HOST, port: process.env.DRACHTIO_PORT, secret: process.env.DRACHTIO_SECRET });
|
||||
srf.on('connect', (err, hp) => {
|
||||
if (err) return this.logger.error({err}, 'Error connecting to drachtio server');
|
||||
logger.info(`connected to drachtio listening on ${hp}`);
|
||||
if (process.env.SBC_ACCOUNT_SID) return;
|
||||
|
||||
const hostports = hp.split(',');
|
||||
for (const hp of hostports) {
|
||||
const arr = /^(.*)\/(.*):(\d+)$/.exec(hp);
|
||||
if (arr && 'udp' === arr[1] && !matcher.contains(arr[2])) {
|
||||
logger.info(`adding sbc address ${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();
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
else {
|
||||
logger.info(`listening in outbound mode on port ${process.env.DRACHTIO_PORT}`);
|
||||
srf.listen({port: process.env.DRACHTIO_PORT, secret: process.env.DRACHTIO_SECRET});
|
||||
}
|
||||
if (process.env.NODE_ENV === 'test') {
|
||||
@@ -140,33 +171,58 @@ 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) {
|
||||
const PORT = process.env.HTTP_PORT || 3000;
|
||||
const getCount = () => activeCallIds.size;
|
||||
const healthCheck = require('@jambonz/http-health-check');
|
||||
healthCheck({port: PORT, logger, path: '/', fn: getCount});
|
||||
}
|
||||
if ('test' !== process.env.NODE_ENV) {
|
||||
/* update call stats periodically */
|
||||
setInterval(() => {
|
||||
stats.gauge('sbc.sip.calls.count', activeCallIds.size, ['direction:inbound']);
|
||||
}, 20000);
|
||||
}
|
||||
|
||||
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 lookupRtpServiceEndpoints = (lookup, serviceName) => {
|
||||
logger.debug(`dns lookup for ${serviceName}..`);
|
||||
lookup(serviceName, {family: 4, all: true}, (err, addresses) => {
|
||||
if (err) {
|
||||
logger.error({err}, `Error looking up ${serviceName}`);
|
||||
return;
|
||||
}
|
||||
logger.debug({addresses, rtpServers}, `dns lookup for ${serviceName} returned`);
|
||||
const addrs = addresses.map((a) => a.address);
|
||||
if (!equalsIgnoreOrder(addrs, rtpServers)) {
|
||||
rtpServers.length = 0;
|
||||
Array.prototype.push.apply(rtpServers, addrs);
|
||||
logger.info({rtpServers}, 'rtpserver endpoints have been updated');
|
||||
setRtpEngines(rtpServers.map((a) => `${a}:${process.env.RTPENGINE_PORT || 22222}`));
|
||||
}
|
||||
});
|
||||
};
|
||||
|
||||
/* 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}`));
|
||||
rtpServers.length = 0;
|
||||
@@ -176,11 +232,31 @@ else {
|
||||
logger.error({err}, 'Error setting new rtpengines');
|
||||
}
|
||||
};
|
||||
|
||||
setInterval(() => {
|
||||
getActiveRtpServers();
|
||||
}, 30000);
|
||||
getActiveRtpServers();
|
||||
|
||||
}
|
||||
|
||||
const {lifecycleEmitter} = require('./lib/autoscale-manager')(logger);
|
||||
|
||||
/* if we are scaling in, check every so often if call count has gone to zero */
|
||||
setInterval(async() => {
|
||||
if (lifecycleEmitter.operationalState === LifeCycleEvents.ScaleIn) {
|
||||
if (0 === activeCallIds.size) {
|
||||
logger.info('scale-in complete now that calls have dried up');
|
||||
lifecycleEmitter.scaleIn();
|
||||
}
|
||||
}
|
||||
}, 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"]
|
||||
}
|
||||
@@ -3,14 +3,14 @@
|
||||
"transport-protocol": "UDP/TLS/RTP/SAVPF",
|
||||
"ICE": "force",
|
||||
"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",
|
||||
"SDES": "off",
|
||||
"flags": ["generate mid", "SDES-no", "media handover"],
|
||||
"flags": ["generate mid", "SDES-no", "media handover", "port latching"],
|
||||
"rtcp-mux": ["accept"]
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,63 @@
|
||||
const noopLogger = {info: () => {}, error: () => {}};
|
||||
const {LifeCycleEvents} = require('./constants');
|
||||
const Emitter = require('events');
|
||||
|
||||
module.exports = (logger) => {
|
||||
logger = logger || noopLogger;
|
||||
|
||||
// listen for SNS lifecycle changes
|
||||
let lifecycleEmitter = new Emitter();
|
||||
lifecycleEmitter.dryUpCalls = false;
|
||||
if (process.env.AWS_SNS_TOPIC_ARM) {
|
||||
|
||||
(async function() {
|
||||
try {
|
||||
lifecycleEmitter = await require('./aws-sns-lifecycle')(logger);
|
||||
|
||||
lifecycleEmitter
|
||||
.on(LifeCycleEvents.ScaleIn, async() => {
|
||||
logger.info('AWS scale-in notification: begin drying up calls');
|
||||
lifecycleEmitter.dryUpCalls = true;
|
||||
lifecycleEmitter.operationalState = LifeCycleEvents.ScaleIn;
|
||||
|
||||
const {srf} = require('..');
|
||||
const {activeCallIds, removeFromRedis} = srf.locals;
|
||||
|
||||
/* remove our private IP from the set of active SBCs so rtp and fs know we are gone */
|
||||
removeFromRedis();
|
||||
|
||||
/* if we have zero calls, we can complete the scale-in right now */
|
||||
const calls = activeCallIds.size;
|
||||
if (0 === calls) {
|
||||
logger.info('scale-in can complete immediately as we have no calls in progress');
|
||||
lifecycleEmitter.completeScaleIn();
|
||||
}
|
||||
else {
|
||||
logger.info(`${calls} calls in progress; scale-in will complete when they are done`);
|
||||
}
|
||||
})
|
||||
.on(LifeCycleEvents.StandbyEnter, () => {
|
||||
lifecycleEmitter.dryUpCalls = true;
|
||||
const {srf} = require('..');
|
||||
const {removeFromRedis} = srf.locals;
|
||||
removeFromRedis();
|
||||
|
||||
logger.info('AWS enter pending state notification: begin drying up calls');
|
||||
})
|
||||
.on(LifeCycleEvents.StandbyExit, () => {
|
||||
lifecycleEmitter.dryUpCalls = false;
|
||||
const {srf} = require('..');
|
||||
const {addToRedis} = srf.locals;
|
||||
addToRedis();
|
||||
|
||||
logger.info('AWS exit pending state notification: re-enable calls');
|
||||
});
|
||||
} catch (err) {
|
||||
logger.error({err}, 'Failure creating SNS notifier, lifecycle events will be disabled');
|
||||
}
|
||||
})();
|
||||
}
|
||||
|
||||
return {lifecycleEmitter};
|
||||
};
|
||||
|
||||
@@ -0,0 +1,186 @@
|
||||
const Emitter = require('events');
|
||||
const bent = require('bent');
|
||||
const assert = require('assert');
|
||||
const PORT = process.env.AWS_SNS_PORT || 3001;
|
||||
const {LifeCycleEvents} = require('./constants');
|
||||
const express = require('express');
|
||||
const app = express();
|
||||
const getString = bent('string');
|
||||
const AWS = require('aws-sdk');
|
||||
const sns = new AWS.SNS({apiVersion: '2010-03-31'});
|
||||
const autoscaling = new AWS.AutoScaling({apiVersion: '2011-01-01'});
|
||||
const {Parser} = require('xml2js');
|
||||
const parser = new Parser();
|
||||
const {validatePayload} = require('verify-aws-sns-signature');
|
||||
|
||||
AWS.config.update({region: process.env.AWS_REGION});
|
||||
|
||||
class SnsNotifier extends Emitter {
|
||||
constructor(logger) {
|
||||
super();
|
||||
|
||||
this.logger = logger;
|
||||
}
|
||||
|
||||
async _handlePost(req, res) {
|
||||
try {
|
||||
const parsedBody = JSON.parse(req.body);
|
||||
this.logger.debug({headers: req.headers, body: parsedBody}, 'Received HTTP POST from AWS');
|
||||
if (!validatePayload(parsedBody)) {
|
||||
this.logger.info('incoming AWS SNS HTTP POST failed signature validation');
|
||||
return res.sendStatus(403);
|
||||
}
|
||||
this.logger.debug('incoming HTTP POST passed validation');
|
||||
res.sendStatus(200);
|
||||
|
||||
switch (parsedBody.Type) {
|
||||
case 'SubscriptionConfirmation':
|
||||
const response = await getString(parsedBody.SubscribeURL);
|
||||
const result = await parser.parseStringPromise(response);
|
||||
this.subscriptionArn = result.ConfirmSubscriptionResponse.ConfirmSubscriptionResult[0].SubscriptionArn[0];
|
||||
this.subscriptionRequestId = result.ConfirmSubscriptionResponse.ResponseMetadata[0].RequestId[0];
|
||||
this.logger.info({
|
||||
subscriptionArn: this.subscriptionArn,
|
||||
subscriptionRequestId: this.subscriptionRequestId
|
||||
}, 'response from SNS SubscribeURL');
|
||||
const data = await this.describeInstance();
|
||||
this.lifecycleState = data.AutoScalingInstances[0].LifecycleState;
|
||||
break;
|
||||
|
||||
case 'Notification':
|
||||
if (parsedBody.Subject.startsWith('Auto Scaling: Lifecycle action \'TERMINATING\'')) {
|
||||
const msg = JSON.parse(parsedBody.Message);
|
||||
if (msg.EC2InstanceId === this.instanceId) {
|
||||
this.logger.info('SnsNotifier - begin scale-in operation');
|
||||
this.scaleInParams = {
|
||||
AutoScalingGroupName: msg.AutoScalingGroupName,
|
||||
LifecycleActionResult: 'CONTINUE',
|
||||
LifecycleActionToken: msg.LifecycleActionToken,
|
||||
LifecycleHookName: msg.LifecycleHookName
|
||||
};
|
||||
this.operationalState = LifeCycleEvents.ScaleIn;
|
||||
this.emit(LifeCycleEvents.ScaleIn);
|
||||
this.unsubscribe();
|
||||
}
|
||||
else {
|
||||
this.logger.debug(`SnsNotifier - instance ${msg.EC2InstanceId} is scaling in (not us)`);
|
||||
}
|
||||
}
|
||||
break;
|
||||
|
||||
default:
|
||||
this.logger.info(`unhandled SNS Post Type: ${parsedBody.Type}`);
|
||||
}
|
||||
|
||||
} catch (err) {
|
||||
this.logger.error({err}, 'Error processing SNS POST request');
|
||||
if (!res.headersSent) res.sendStatus(500);
|
||||
}
|
||||
}
|
||||
|
||||
async init() {
|
||||
try {
|
||||
this.logger.info('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
|
||||
}, 'retrieved AWS instance data');
|
||||
|
||||
// start listening
|
||||
app.use(express.urlencoded({ extended: true }));
|
||||
app.use(express.json());
|
||||
app.use(express.text());
|
||||
app.post('/', this._handlePost.bind(this));
|
||||
app.use((err, req, res, next) => {
|
||||
this.logger.error(err, 'burped error');
|
||||
res.status(err.status || 500).json({msg: err.message});
|
||||
});
|
||||
app.listen(PORT);
|
||||
|
||||
} catch (err) {
|
||||
this.logger.error({err}, 'Error retrieving AWS instance metadata');
|
||||
}
|
||||
}
|
||||
|
||||
async subscribe() {
|
||||
try {
|
||||
const response = await sns.subscribe({
|
||||
Protocol: 'http',
|
||||
TopicArn: process.env.AWS_SNS_TOPIC_ARM,
|
||||
Endpoint: this.snsEndpoint
|
||||
}).promise();
|
||||
this.logger.info({response}, `response to SNS subscribe to ${process.env.AWS_SNS_TOPIC_ARM}`);
|
||||
} catch (err) {
|
||||
this.logger.error({err}, `Error subscribing to SNS topic arn ${process.env.AWS_SNS_TOPIC_ARM}`);
|
||||
}
|
||||
}
|
||||
|
||||
async unsubscribe() {
|
||||
if (!this.subscriptionArn) throw new Error('SnsNotifier#unsubscribe called without an active subscription');
|
||||
try {
|
||||
const response = await sns.unsubscribe({
|
||||
SubscriptionArn: this.subscriptionArn
|
||||
}).promise();
|
||||
this.logger.info({response}, `response to SNS unsubscribe to ${process.env.AWS_SNS_TOPIC_ARM}`);
|
||||
} catch (err) {
|
||||
this.logger.error({err}, `Error unsubscribing to SNS topic arn ${process.env.AWS_SNS_TOPIC_ARM}`);
|
||||
}
|
||||
}
|
||||
|
||||
completeScaleIn() {
|
||||
assert(this.scaleInParams);
|
||||
autoscaling.completeLifecycleAction(this.scaleInParams, (err, response) => {
|
||||
if (err) return this.logger.error({err}, 'Error completing scale-in');
|
||||
this.logger.info({response}, 'Successfully completed scale-in action');
|
||||
});
|
||||
}
|
||||
|
||||
describeInstance() {
|
||||
return new Promise((resolve, reject) => {
|
||||
if (!this.instanceId) return reject('instance-id unknown');
|
||||
autoscaling.describeAutoScalingInstances({
|
||||
InstanceIds: [this.instanceId]
|
||||
}, (err, data) => {
|
||||
if (err) {
|
||||
this.logger.error({err}, 'Error describing instances');
|
||||
reject(err);
|
||||
} else {
|
||||
this.logger.info({data}, 'SnsNotifier: describeInstance');
|
||||
resolve(data);
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
module.exports = async function(logger) {
|
||||
const notifier = new SnsNotifier(logger);
|
||||
await notifier.init();
|
||||
await notifier.subscribe();
|
||||
|
||||
process.on('SIGHUP', async() => {
|
||||
try {
|
||||
const data = await notifier.describeInstance();
|
||||
const state = data.AutoScalingInstances[0].LifecycleState;
|
||||
if (state !== notifier.lifecycleState) {
|
||||
notifier.lifecycleState = state;
|
||||
switch (state) {
|
||||
case 'Standby':
|
||||
notifier.emit(LifeCycleEvents.StandbyEnter);
|
||||
break;
|
||||
case 'InService':
|
||||
notifier.emit(LifeCycleEvents.StandbyExit);
|
||||
break;
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
console.error(err);
|
||||
}
|
||||
});
|
||||
return notifier;
|
||||
};
|
||||
+188
-21
@@ -1,7 +1,7 @@
|
||||
const Emitter = require('events');
|
||||
const {makeRtpEngineOpts, SdpWantsSrtp, makeCallCountKey} = 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';
|
||||
@@ -13,8 +13,8 @@ const MS_TEAMS_SIP_ENDPOINT = 'sip.pstnhub.microsoft.com';
|
||||
const createBLegFromHeader = (req) => {
|
||||
const from = req.getParsedHeader('From');
|
||||
const uri = parseUri(from.uri);
|
||||
if (uri && uri.user) return `sip:${uri.user}@localhost`;
|
||||
return 'sip:anonymous@localhost';
|
||||
if (uri && uri.user) return `<sip:${uri.user}@localhost>`;
|
||||
return '<sip:anonymous@localhost>';
|
||||
};
|
||||
|
||||
class CallSession extends Emitter {
|
||||
@@ -39,6 +39,10 @@ class CallSession extends Emitter {
|
||||
return !!this.req.locals.msTeamsTenantFqdn;
|
||||
}
|
||||
|
||||
get privateSipAddress() {
|
||||
return this.srf.locals.privateSipAddress;
|
||||
}
|
||||
|
||||
async connect() {
|
||||
this.logger.info('inbound call accepted for routing');
|
||||
const engine = this.getRtpEngine();
|
||||
@@ -49,10 +53,26 @@ 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,
|
||||
subscribeDTMF,
|
||||
unsubscribeDTMF
|
||||
} = engine;
|
||||
this.offer = offer;
|
||||
this.answer = answer;
|
||||
this.del = del;
|
||||
this.blockMedia = blockMedia;
|
||||
this.unblockMedia = unblockMedia;
|
||||
this.blockDTMF = blockDTMF;
|
||||
this.unblockDTMF = unblockDTMF;
|
||||
this.subscribeDTMF = subscribeDTMF;
|
||||
this.unsubscribeDTMF = unsubscribeDTMF;
|
||||
|
||||
const featureServer = await this.getFeatureServer();
|
||||
if (!featureServer) {
|
||||
@@ -98,15 +118,22 @@ class CallSession extends Emitter {
|
||||
}
|
||||
|
||||
// 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 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});
|
||||
|
||||
@@ -136,7 +163,14 @@ class CallSession extends Emitter {
|
||||
proxy,
|
||||
headers,
|
||||
responseHeaders,
|
||||
proxyRequestHeaders: ['all', '-Authorization', '-Max-Forwards', '-Record-Route', '-Session-Expires', 'Min-SE'],
|
||||
proxyRequestHeaders: [
|
||||
'all',
|
||||
'-Authorization',
|
||||
'-Max-Forwards',
|
||||
'-Record-Route',
|
||||
'-Session-Expires',
|
||||
'-X-Subspace-Forwarded-For'
|
||||
],
|
||||
proxyResponseHeaders: ['all'],
|
||||
localSdpB: response.sdp,
|
||||
localSdpA: async(sdp, res) => {
|
||||
@@ -163,7 +197,7 @@ 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) {
|
||||
@@ -175,17 +209,22 @@ class CallSession extends Emitter {
|
||||
else if (err.message !== 'call canceled') {
|
||||
this.logger.error(err, 'unexpected error routing inbound call');
|
||||
}
|
||||
this.srf.endSession(this.req);
|
||||
}
|
||||
}
|
||||
|
||||
_setDlgHandlers(dlg) {
|
||||
this.activeCallIds.set(this.req.get('Call-ID'), this);
|
||||
const {callId} = dlg.sip;
|
||||
this.activeCallIds.set(callId, this);
|
||||
this.subscribeDTMF(this.logger, callId, this.rtpEngineOpts.uas.tag,
|
||||
this._onDTMF.bind(this));
|
||||
dlg.on('destroy', () => {
|
||||
debug('call ended with normal termination');
|
||||
this.logger.info('call ended with normal termination');
|
||||
this.rtpEngineResource.destroy().catch((err) => {});
|
||||
this.activeCallIds.delete(this.req.get('Call-ID'));
|
||||
this.activeCallIds.delete(callId);
|
||||
if (dlg.other && dlg.other.connected) dlg.other.destroy().catch((e) => {});
|
||||
this.srf.endSession(this.req);
|
||||
});
|
||||
|
||||
//re-invite
|
||||
@@ -213,7 +252,7 @@ class CallSession extends Emitter {
|
||||
this.rtpEngineResource.destroy().catch((err) => {});
|
||||
this.activeCallIds.delete(this.req.get('Call-ID'));
|
||||
dlg.other.destroy().catch((e) => {});
|
||||
|
||||
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
|
||||
if (process.env.JAMBONES_HOSTING) {
|
||||
this.decrKey(this.callCountKey)
|
||||
.then((count) => this.logger.debug({key: this.callCountKey},
|
||||
@@ -235,16 +274,52 @@ class CallSession extends Emitter {
|
||||
trunk
|
||||
}).catch((err) => this.logger.error({err}, 'Error writing cdr for completed call'));
|
||||
}
|
||||
this.srf.endSession(this.req);
|
||||
});
|
||||
});
|
||||
|
||||
this.subscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag,
|
||||
this._onDTMF.bind(this, uac));
|
||||
|
||||
uas.on('modify', this._onReinvite.bind(this, uas));
|
||||
uac.on('modify', this._onReinvite.bind(this, uac));
|
||||
|
||||
uac.on('refer', this._onFeatureServerTransfer.bind(this, uac));
|
||||
uas.on('refer', this._onRefer.bind(this, uas));
|
||||
|
||||
uas.on('info', this._onInfo.bind(this, uas));
|
||||
uac.on('info', this._onInfo.bind(this, uac));
|
||||
|
||||
// default forwarding of other request types
|
||||
forwardInDialogRequests(uas, ['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 +329,7 @@ class CallSession extends Emitter {
|
||||
*/
|
||||
async replaces(req, res) {
|
||||
try {
|
||||
let opts = Object.assign(this.rtpEngineOpts.offer, {sdp: req.body});
|
||||
let opts = Object.assign({}, this.rtpEngineOpts.uas.mediaOpts, {sdp: req.body});
|
||||
let response = await this.offer(opts);
|
||||
if ('ok' !== response.result) {
|
||||
res.send(488);
|
||||
@@ -262,8 +337,8 @@ 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 = Object.assign({}, this.rtpEngineOpts.uac.mediaOpts, {sdp, 'to-tag': this.toTag});
|
||||
Object.assign(this.rtpEngineOpts.uas.mediaOpts, {'to-tag': this.toTag});
|
||||
response = await this.answer(opts);
|
||||
if ('ok' !== response.result) {
|
||||
res.send(488);
|
||||
@@ -281,6 +356,8 @@ class CallSession extends Emitter {
|
||||
});
|
||||
}
|
||||
|
||||
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
|
||||
|
||||
const uas = await this.srf.createUAS(req, res, {
|
||||
localSdp: response.sdp,
|
||||
headers
|
||||
@@ -301,11 +378,18 @@ class CallSession extends Emitter {
|
||||
|
||||
async _onReinvite(dlg, req, res) {
|
||||
try {
|
||||
/* check for re-invite with no SDP -- seen that from BT when they provide UUI info */
|
||||
if (!req.body) {
|
||||
this.logger.info('got a reINVITE with no SDP; just respond with our current offer');
|
||||
res.send(200, {body: dlg.local.sdp});
|
||||
return;
|
||||
}
|
||||
const reason = req.get('X-Reason');
|
||||
const fromTag = dlg.type === 'uas' ? this.rtpEngineOpts.uas.tag : this.rtpEngineOpts.uac.tag;
|
||||
const toTag = dlg.type === 'uas' ? this.rtpEngineOpts.uac.tag : this.rtpEngineOpts.uas.tag;
|
||||
const offerMedia = dlg.type === 'uas' ? this.rtpEngineOpts.uac.mediaOpts : this.rtpEngineOpts.uas.mediaOpts;
|
||||
const answerMedia = dlg.type === 'uas' ? this.rtpEngineOpts.uas.mediaOpts : this.rtpEngineOpts.uac.mediaOpts;
|
||||
const direction = dlg.type === 'uas' ? ['private', 'public'] : ['public', 'private'];
|
||||
const direction = dlg.type === 'uas' ? ['public', 'private'] : ['private', 'public'];
|
||||
let opts = {
|
||||
...this.rtpEngineOpts.common,
|
||||
...offerMedia,
|
||||
@@ -314,13 +398,23 @@ class CallSession extends Emitter {
|
||||
direction,
|
||||
sdp: req.body,
|
||||
};
|
||||
//if (reason && opts.flags && !opts.flags.includes('reset')) opts.flags.push('reset');
|
||||
|
||||
let response = await this.offer(opts);
|
||||
if ('ok' !== response.result) {
|
||||
res.send(488);
|
||||
throw new Error(`_onReinvite: rtpengine failed: offer: ${JSON.stringify(response)}`);
|
||||
}
|
||||
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 (reason && dlg.type === 'uac' && ['release-media', 'anchor-media'].includes(reason)) {
|
||||
this.logger.info({response}, `got a reinvite from FS to ${reason}`);
|
||||
sdp = dlg.other.remote.sdp;
|
||||
}
|
||||
else {
|
||||
sdp = await dlg.other.modify(response.sdp);
|
||||
}
|
||||
opts = {
|
||||
...this.rtpEngineOpts.common,
|
||||
...answerMedia,
|
||||
@@ -339,15 +433,86 @@ class CallSession extends Emitter {
|
||||
}
|
||||
}
|
||||
|
||||
async _onInfo(dlg, req, res) {
|
||||
const fromTag = dlg.type === 'uas' ? this.rtpEngineOpts.uas.tag : this.rtpEngineOpts.uac.tag;
|
||||
try {
|
||||
if (dlg.type === 'uac' && req.has('X-Reason')) {
|
||||
const reason = req.get('X-Reason');
|
||||
const opts = {
|
||||
...this.rtpEngineOpts.common,
|
||||
flags: ['reset'],
|
||||
'from-tag': fromTag
|
||||
};
|
||||
this.logger.info(`_onInfo: got request ${reason}`);
|
||||
res.send(200);
|
||||
|
||||
if (reason.startsWith('mute')) {
|
||||
const response = Promise.all([this.blockMedia(opts), this.blockDTMF(opts)]);
|
||||
this.logger.info({response}, `_onInfo: response to rtpengine command for ${reason}`);
|
||||
}
|
||||
else if (reason.startsWith('unmute')) {
|
||||
const response = Promise.all([this.unblockMedia(opts), this.unblockDTMF(opts)]);
|
||||
this.logger.info({response}, `_onInfo: response to rtpengine command for ${reason}`);
|
||||
}
|
||||
}
|
||||
else {
|
||||
const immutableHdrs = ['via', 'from', 'to', 'call-id', 'cseq', 'max-forwards', 'content-length'];
|
||||
const headers = {};
|
||||
Object.keys(req.headers).forEach((h) => {
|
||||
if (!immutableHdrs.includes(h)) headers[h] = req.headers[h];
|
||||
});
|
||||
const response = await dlg.other.request({method: 'INFO', headers, body: req.body});
|
||||
const responseHeaders = {};
|
||||
if (response.has('Content-Type')) {
|
||||
Object.assign(responseHeaders, {'Content-Type': response.get('Content-Type')});
|
||||
}
|
||||
res.send(response.status, {headers: responseHeaders, body: response.body});
|
||||
}
|
||||
} catch (err) {
|
||||
this.logger.info({err}, `Error handing INFO request on ${dlg.type} leg`);
|
||||
}
|
||||
}
|
||||
|
||||
async _onFeatureServerTransfer(dlg, req, res) {
|
||||
try {
|
||||
const referTo = req.getParsedHeader('Refer-To');
|
||||
const uri = parseUri(referTo.uri);
|
||||
this.logger.info({uri, referTo}, '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);
|
||||
|
||||
let selectedGateway = false;
|
||||
let e164 = false;
|
||||
if (gateway) {
|
||||
/* host of Refer-to to an outbound gateway */
|
||||
const gw = await this.srf.locals.getOutboundGatewayForRefer(gateway.voip_carrier_sid);
|
||||
if (gw) {
|
||||
selectedGateway = true;
|
||||
e164 = gw.e164_leading_plus;
|
||||
uri.host = gw.ipv4;
|
||||
uri.port = gw.port;
|
||||
}
|
||||
}
|
||||
if (!selectedGateway) {
|
||||
uri.host = this.req.source_address;
|
||||
uri.port = this.req.source_port;
|
||||
}
|
||||
if (e164 && !uri.user.startsWith('+')) {
|
||||
uri.user = `+${uri.user}`;
|
||||
}
|
||||
const response = await this.uas.request({
|
||||
method: 'REFER',
|
||||
headers: {
|
||||
'Refer-To': stringifyUri(uri),
|
||||
'Referred-By': stringifyUri(u)
|
||||
}
|
||||
});
|
||||
return res.send(response.status);
|
||||
}
|
||||
res.send(202);
|
||||
|
||||
@@ -367,6 +532,7 @@ class CallSession extends Emitter {
|
||||
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(() => {});
|
||||
@@ -438,6 +604,7 @@ class CallSession extends Emitter {
|
||||
|
||||
// successfully connected
|
||||
this.logger.info('successfully connected new call leg for REFER');
|
||||
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
|
||||
this.referInvite = null;
|
||||
sendNotify(this.uas, '200 OK');
|
||||
this.uas.destroy();
|
||||
|
||||
@@ -0,0 +1,7 @@
|
||||
{
|
||||
"LifeCycleEvents" : {
|
||||
"ScaleIn": "scale-in",
|
||||
"StandbyEnter": "standby-enter",
|
||||
"StandbyExit": "standby-exit"
|
||||
}
|
||||
}
|
||||
+56
-5
@@ -11,6 +11,15 @@ 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
|
||||
@@ -35,6 +44,13 @@ SELECT * FROM phone_numbers
|
||||
WHERE number = ?
|
||||
AND voip_carrier_sid = ?`;
|
||||
|
||||
const sqlSelectOutboundGatewayForCarrier = `
|
||||
SELECT ipv4, port, e164_leading_plus
|
||||
FROM sip_gateways sg, voip_carriers vc
|
||||
WHERE sg.voip_carrier_sid = ?
|
||||
AND sg.voip_carrier_sid = vc.voip_carrier_sid
|
||||
AND outbound = 1`;
|
||||
|
||||
const gatewayMatchesSourceAddress = (source_address, gw) => {
|
||||
if (32 === gw.netmask && gw.ipv4 === source_address) return true;
|
||||
if (gw.netmask < 32) {
|
||||
@@ -48,6 +64,19 @@ module.exports = (srf, logger) => {
|
||||
const {pool} = srf.locals.dbHelpers;
|
||||
const pp = pool.promise();
|
||||
|
||||
const getOutboundGatewayForRefer = async(voip_carrier_sid) => {
|
||||
try {
|
||||
const [r] = await pp.query(sqlSelectOutboundGatewayForCarrier, [voip_carrier_sid]);
|
||||
if (0 === r.length) return null;
|
||||
|
||||
/* if multiple, prefer a DNS name */
|
||||
const hasDns = r.find((row) => row.ipv4.match(/^[A-Za-z]/));
|
||||
return hasDns || r[0];
|
||||
} catch (err) {
|
||||
logger.error({err}, 'getOutboundGatewayForRefer');
|
||||
}
|
||||
};
|
||||
|
||||
const getApplicationForDidAndCarrier = async(req, voip_carrier_sid) => {
|
||||
const did = normalizeDID(req.calledNumber);
|
||||
|
||||
@@ -73,7 +102,11 @@ module.exports = (srf, logger) => {
|
||||
const [r] = await pp.query(sqlCarriersForAccountBySid,
|
||||
[process.env.SBC_ACCOUNT_SID, req.source_address, req.source_port]);
|
||||
if (0 === r.length) return failure;
|
||||
return {fromCarrier: true, gateway: r[0], account_sid: process.env.SBC_ACCOUNT_SID};
|
||||
return {
|
||||
fromCarrier: true,
|
||||
gateway: r[0],
|
||||
account_sid: process.env.SBC_ACCOUNT_SID
|
||||
};
|
||||
}
|
||||
else {
|
||||
/* we may have a carrier at the service provider level */
|
||||
@@ -88,8 +121,23 @@ module.exports = (srf, logger) => {
|
||||
|
||||
const [r] = await pp.query(sql);
|
||||
if (0 === r.length) {
|
||||
/* came from a carrier, but number is not provisioned */
|
||||
return {fromCarrier: true};
|
||||
/* came from a carrier, but number is not provisioned..
|
||||
check if we only have a single account, otherwise we have no
|
||||
way of knowing which account this is for
|
||||
*/
|
||||
const [r] = await pp.query('SELECT count(*) as count from accounts where service_provider_sid = ?',
|
||||
matches[0].service_provider_sid);
|
||||
if (r[0].count === 0) return {fromCarrier: true};
|
||||
else {
|
||||
const [accounts] = await pp.query('SELECT * from accounts where service_provider_sid = ?',
|
||||
matches[0].service_provider_sid);
|
||||
return {
|
||||
fromCarrier: true,
|
||||
gateway: matches[0],
|
||||
account_sid: accounts[0].account_sid,
|
||||
account: accounts[0]
|
||||
};
|
||||
}
|
||||
}
|
||||
const gateway = matches.find((m) => m.voip_carrier_sid === r[0].voip_carrier_sid);
|
||||
const [accounts] = await pp.query(sqlAccountBySid, r[0].account_sid);
|
||||
@@ -107,7 +155,9 @@ module.exports = (srf, logger) => {
|
||||
}
|
||||
|
||||
/* get all the carriers and gateways for the account owning this sip realm */
|
||||
const [gw] = await pp.query(sqlSelectAllCarriersForAccountByRealm, uri.host);
|
||||
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);
|
||||
@@ -125,6 +175,7 @@ module.exports = (srf, logger) => {
|
||||
|
||||
return {
|
||||
wasOriginatedFromCarrier,
|
||||
getApplicationForDidAndCarrier
|
||||
getApplicationForDidAndCarrier,
|
||||
getOutboundGatewayForRefer
|
||||
};
|
||||
};
|
||||
|
||||
+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}`);
|
||||
}
|
||||
|
||||
+31
-8
@@ -68,11 +68,20 @@ module.exports = function(srf, logger) {
|
||||
blacklistUnknownRealms: true,
|
||||
emitter: new AuthOutcomeReporter(stats)
|
||||
});
|
||||
const {wasOriginatedFromCarrier, getApplicationForDidAndCarrier} = require('./db-utils')(srf, logger);
|
||||
|
||||
|
||||
const initLocals = (req, res, next) => {
|
||||
req.locals = req.locals || {};
|
||||
|
||||
/* check if forwarded by a proxy that applied an X-Forwarded-For Header */
|
||||
if (req.has('X-Forwarded-For') || req.has('X-Subspace-Forwarded-For')) {
|
||||
const original_source_address = req.get('X-Forwarded-For') || req.get('X-Subspace-Forwarded-For');
|
||||
logger.info({
|
||||
callId: req.get('Call-ID'),
|
||||
original_source_address,
|
||||
proxy_source_address: req.source_address,
|
||||
}, 'overwriting source address for proxied SIP INVITE');
|
||||
req.source_address = original_source_address;
|
||||
}
|
||||
req.locals.cdr = initCdr(req);
|
||||
const callId = req.get('Call-ID');
|
||||
req.on('cancel', () => {
|
||||
@@ -102,8 +111,14 @@ module.exports = function(srf, logger) {
|
||||
|
||||
const identifyAccount = async(req, res, next) => {
|
||||
try {
|
||||
|
||||
const {fromCarrier, gateway, account_sid, application_sid, account} = await wasOriginatedFromCarrier(req);
|
||||
const {wasOriginatedFromCarrier, getApplicationForDidAndCarrier} = req.srf.locals;
|
||||
const {
|
||||
fromCarrier,
|
||||
gateway,
|
||||
account_sid,
|
||||
application_sid,
|
||||
account
|
||||
} = await wasOriginatedFromCarrier(req);
|
||||
/**
|
||||
* calls come from 3 sources:
|
||||
* (1) A carrier
|
||||
@@ -122,6 +137,8 @@ module.exports = function(srf, logger) {
|
||||
req.locals = {
|
||||
originator: 'trunk',
|
||||
carrier: gateway.name,
|
||||
gateway,
|
||||
voip_carrier_sid: gateway.voip_carrier_sid,
|
||||
application_sid: sid || gateway.application_sid,
|
||||
account_sid,
|
||||
account,
|
||||
@@ -135,7 +152,8 @@ module.exports = function(srf, logger) {
|
||||
const app = await lookupAppByTeamsTenant(uri.host);
|
||||
if (!app) {
|
||||
stats.increment('sbc.terminations', ['sipStatus:404']);
|
||||
return res.send(404, {headers: {'X-Reason': 'no configured application'}});
|
||||
res.send(404, {headers: {'X-Reason': 'no configured application'}});
|
||||
return req.srf.endSession(req);
|
||||
}
|
||||
|
||||
req.locals = {
|
||||
@@ -154,7 +172,8 @@ module.exports = function(srf, logger) {
|
||||
const account = await lookupAccountBySipRealm(uri.host);
|
||||
if (!account) {
|
||||
stats.increment('sbc.terminations', ['sipStatus:404']);
|
||||
return res.send(404);
|
||||
res.send(404);
|
||||
return req.srf.endSession(req);
|
||||
}
|
||||
|
||||
/* if this is a dedicated SBC (static IP) only take calls for that account's sip realm */
|
||||
@@ -163,7 +182,8 @@ module.exports = function(srf, logger) {
|
||||
`identifyAccount: static IP for ${process.env.SBC_ACCOUNT_SID} but call for ${account.account_sid}`);
|
||||
stats.increment('sbc.terminations', ['sipStatus:404']);
|
||||
delete req.locals.cdr;
|
||||
return res.send(404);
|
||||
res.send(404);
|
||||
return req.srf.endSession(req);
|
||||
}
|
||||
req.locals = {
|
||||
account_sid: account.account_sid,
|
||||
@@ -244,13 +264,15 @@ module.exports = function(srf, logger) {
|
||||
account_sid,
|
||||
count: limit_sessions
|
||||
}).catch((err) => logger.info({err}, 'checkLimits: error writing alert'));
|
||||
return res.send(503, 'Maximum Calls In Progress');
|
||||
res.send(503, 'Maximum Calls In Progress');
|
||||
return req.srf.endSession(req);
|
||||
}
|
||||
next();
|
||||
} catch (err) {
|
||||
stats.increment('sbc.terminations', ['sipStatus:500']);
|
||||
logger.error({err}, 'error checking limits error for inbound call');
|
||||
res.send(500);
|
||||
req.srf.endSession(req);
|
||||
}
|
||||
};
|
||||
|
||||
@@ -263,6 +285,7 @@ module.exports = function(srf, logger) {
|
||||
stats.increment('sbc.terminations', ['sipStatus:500']);
|
||||
logger.error(err, `${req.get('Call-ID')} Error looking up related info for inbound call`);
|
||||
res.send(500);
|
||||
req.srf.endSession(req);
|
||||
}
|
||||
};
|
||||
|
||||
|
||||
+13
-1
@@ -45,11 +45,23 @@ const normalizeDID = (tel) => {
|
||||
return arr ? arr[1] : tel;
|
||||
};
|
||||
|
||||
const equalsIgnoreOrder = (a, b) => {
|
||||
if (a.length !== b.length) return false;
|
||||
const uniqueValues = new Set([...a, ...b]);
|
||||
for (const v of uniqueValues) {
|
||||
const aCount = a.filter((e) => e === v).length;
|
||||
const bCount = b.filter((e) => e === v).length;
|
||||
if (aCount !== bCount) return false;
|
||||
}
|
||||
return true;
|
||||
};
|
||||
|
||||
module.exports = {
|
||||
isWSS,
|
||||
SdpWantsSrtp,
|
||||
getAppserver,
|
||||
makeRtpEngineOpts,
|
||||
makeCallCountKey,
|
||||
normalizeDID
|
||||
normalizeDID,
|
||||
equalsIgnoreOrder
|
||||
};
|
||||
|
||||
Generated
+2360
-674
File diff suppressed because it is too large
Load Diff
+20
-14
@@ -1,9 +1,9 @@
|
||||
{
|
||||
"name": "sbc-inbound",
|
||||
"version": "0.3.6",
|
||||
"version": "v0.7.1",
|
||||
"main": "app.js",
|
||||
"engines": {
|
||||
"node": ">= 10.16.0"
|
||||
"node": ">= 12.0.0"
|
||||
},
|
||||
"keywords": [
|
||||
"sip",
|
||||
@@ -20,28 +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 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/db-helpers": "^0.6.16",
|
||||
"@jambonz/http-authenticator": "^0.2.0",
|
||||
"@jambonz/realtimedb-helpers": "^0.4.3",
|
||||
"@jambonz/rtpengine-utils": "^0.1.12",
|
||||
"@jambonz/stats-collector": "^0.1.5",
|
||||
"@jambonz/time-series": "^0.1.5",
|
||||
"@jambonz/http-health-check": "^0.0.1",
|
||||
"@jambonz/realtimedb-helpers": "^0.4.9",
|
||||
"@jambonz/rtpengine-utils": "^0.3.0-beta.3",
|
||||
"@jambonz/stats-collector": "^0.1.6",
|
||||
"@jambonz/time-series": "^0.1.6",
|
||||
"aws-sdk": "^2.1036.0",
|
||||
"bent": "^7.3.12",
|
||||
"cidr-matcher": "^2.1.1",
|
||||
"debug": "^4.3.1",
|
||||
"debug": "^4.3.3",
|
||||
"drachtio-fn-b2b-sugar": "0.0.12",
|
||||
"drachtio-srf": "^4.4.49",
|
||||
"pino": "^6.8.0",
|
||||
"rtpengine-client": "^0.2.0"
|
||||
"drachtio-srf": "^4.4.59",
|
||||
"express": "^4.17.1",
|
||||
"pino": "^7.4.1",
|
||||
"rtpengine-client": "^0.2.0",
|
||||
"verify-aws-sns-signature": "^0.0.6",
|
||||
"xml2js": "^0.4.23"
|
||||
},
|
||||
"devDependencies": {
|
||||
"clear-module": "^4.1.1",
|
||||
"eslint": "^7.15.0",
|
||||
"eslint-plugin-promise": "^4.2.1",
|
||||
"eslint": "^7.32.0",
|
||||
"eslint-plugin-promise": "^6.0.0",
|
||||
"nyc": "^15.1.0",
|
||||
"tape": "^4.13.3"
|
||||
}
|
||||
|
||||
@@ -10,6 +10,7 @@ networks:
|
||||
services:
|
||||
mysql:
|
||||
image: mysql:5.7
|
||||
platform: linux/x86_64
|
||||
ports:
|
||||
- "3306:3306"
|
||||
environment:
|
||||
@@ -71,7 +72,7 @@ services:
|
||||
ipv4_address: 172.38.0.14
|
||||
|
||||
influxdb:
|
||||
image: influxdb:1.8-alpine
|
||||
image: influxdb:1.8
|
||||
ports:
|
||||
- "8086:8086"
|
||||
networks:
|
||||
|
||||
Reference in New Issue
Block a user