Compare commits

...
42 Commits
Author SHA1 Message Date
two56andMatt Preskett b6675de2bd Fix: Feature server REFER (#114)
* Update _onFeatureServerTransfer to fix REFER case

* Add INFO listener

---------

Co-authored-by: Matt Preskett <matt.preskett@netcall.com>
2023-08-20 08:47:01 -04:00
Dave Horton 66e73d8353 logging 2023-08-14 14:36:07 -04:00
Hoan Luu Huu a25aa08e17 update stats colector version (#111) 2023-08-07 21:22:47 -04:00
Hoan Luu Huu b1296c8e89 correct pare callRecording headers from sip info (#110) 2023-07-20 09:02:04 -04:00
Hoan Luu Huu 17df3a4ca8 Feat/siprec custom headers (#109)
* siprec custom headers

* siprec custom headers

* update siprec client util
2023-07-20 08:21:09 -04:00
Hoan Luu Huu f2e820915d Multi srs (#106)
* multi srs

* multi srs

* multi srs

* fix review comment
2023-07-04 16:41:46 +01:00
Dave Horton a1ec0fd5da 0.8.4 2023-06-28 09:33:01 +01:00
Hoan Luu Huu d7654de526 fix: update aws sdk v3 (#108)
* fix: update aws sdk v3

* fix: update aws sdk v3

* fix jslint issue

* fix jslint issue

* fix parse aws response

* fix parse aws response
2023-06-28 09:20:37 +01:00
Hoan Luu Huu 6e47e37cc2 fix client encrypted password (#105)
* fix client encrypted password

* fix client encrypted password

* fix failing testcase
2023-06-15 20:47:31 -04:00
Hoan Luu Huu c37d417a69 Feat: jambonz clients (#104)
* authenticate user

* authenticate user

* update db helper

* update db helper

* fix jslint

* use digest-utils
2023-06-15 07:36:33 -04:00
Hoan Luu Huu f7b3815ee3 feat: record all calls (#98)
* feat: record all calls

* wip: add cdr

* wip: add cdr

* fix

* fix jslint

* fix: record format ext
2023-06-09 14:57:45 -04:00
Hoan Luu Huu f385a33297 redis sentinel configuration (#103)
* redis sentinel configuration

* redis sentinel configuration

* update redis version

* update redis version
2023-06-07 10:00:44 -04:00
Hoan Luu Huu a01febdc22 update dbhelper and redis (#101) 2023-06-01 08:32:32 -04:00
Dave Horton e641c590b2 Fix/cidr error handling (#102)
* fix docker build

* catch error from CIDR which can happen with invalid sip gateway data
2023-05-31 09:11:53 -04:00
Dave Horton 8448e003f6 update siprec-client-utils 2023-05-29 09:23:37 -04:00
Dave Horton f6e071f31e standardize on passing .query args as array (#99) 2023-05-22 09:56:48 -04:00
Dave Horton 54d9044937 0.8.3 2023-05-11 09:25:24 -04:00
Dave Horton bcc089058f update deps 2023-05-08 13:12:50 -04:00
Snyk bot 1a09a13cfa fix: package.json & package-lock.json to reduce vulnerabilities (#97)
The following vulnerabilities are fixed with an upgrade:
- https://snyk.io/vuln/SNYK-JS-XML2JS-5414874
2023-04-25 07:47:16 -04:00
Hoan Luu Huu b311d4af48 fix: ice issue (#95)
* fix: ice issue

* fix: review comment
2023-04-13 08:43:56 -04:00
Anton Voylenko 0fe8843729 README and env variables validation (#91)
* validate env at startup

* README updated with environment variables
2023-04-10 15:36:36 -04:00
Dave Horton 0b0e37020c push to docker 2023-04-10 09:43:00 -04:00
Antony Jukes 01cde07117 if else made authorization check unreachable (#90) 2023-04-06 09:34:37 -04:00
Hoan Luu Huu 31b4534878 fix: update stat collector version (#89) 2023-04-05 12:03:36 -04:00
Hoan Luu Huu 0c6365a8df fix: update stat collector version (#89) 2023-04-05 12:03:17 -04:00
Dave Hortonandsnyk-bot def2cb2be2 fix: Dockerfile to reduce vulnerabilities (#88)
The following vulnerabilities are fixed with an upgrade:
- https://snyk.io/vuln/SNYK-ALPINE316-OPENSSL-3368756
- https://snyk.io/vuln/SNYK-ALPINE316-OPENSSL-3368756
- https://snyk.io/vuln/SNYK-ALPINE316-OPENSSL-5291792
- https://snyk.io/vuln/SNYK-ALPINE316-OPENSSL-5291792

Co-authored-by: snyk-bot <snyk-bot@snyk.io>
2023-03-30 13:31:30 -04:00
Dave Horton baaee6f063 bump version 2023-03-28 14:16:07 -04:00
Hoan Luu HuuandQuan HL d2f32229a2 feat: add instance_id to gause calls metric (#87)
Co-authored-by: Quan HL <quanluuhoang8@gmail.com>
2023-03-21 07:54:52 -04:00
Dave Horton 38f2d27246 refactor of speech-utils (#85) 2023-03-14 10:01:19 -04:00
Dave Horton d94977169e pass on displayName of From header if we get it (#82)
* pass on displayName of From header if we get it

* bump version
2023-03-03 18:55:56 -05:00
Dave Horton 0c9df878b8 update dockerfile 2023-02-23 08:23:17 -05:00
Dave Horton 93930fff18 fix for #79 (#80) 2023-02-22 08:51:35 -05:00
Dave Horton 2ce24659ca initial attempt (#77)
* initial attempt

* fix: special handling on uas dialog
2023-02-21 14:00:36 -05:00
EgleH 7c10e98a66 Upgrade node to node:18.14.0-alpine3.16 (#78) 2023-02-21 07:54:55 -05:00
Dave Horton f3fe46319d update siprec client with fix for stop recording 2023-02-17 12:10:14 -05:00
Dave Horton 43bfeb439a update siprec-client with fix for smart tap xml 2023-02-14 14:15:11 -05:00
Dave Horton 31f34fbffa update to latest siprec client 2023-02-13 16:50:07 -05:00
Dave Horton 52d0670b88 bump version 2023-02-13 09:15:05 -05:00
Snyk bot f42128765d fix: Dockerfile to reduce vulnerabilities (#75)
The following vulnerabilities are fixed with an upgrade:
- https://snyk.io/vuln/SNYK-ALPINE316-OPENSSL-3314623
- https://snyk.io/vuln/SNYK-ALPINE316-OPENSSL-3314624
- https://snyk.io/vuln/SNYK-ALPINE316-OPENSSL-3314624
- https://snyk.io/vuln/SNYK-ALPINE316-OPENSSL-3314641
- https://snyk.io/vuln/SNYK-ALPINE316-OPENSSL-3314643
2023-02-12 12:45:19 -05:00
Dave Horton eac1d7a9e9 K8s scale in (#73)
* changes for proper scale-in for K8S

* fix check for call count
2023-02-07 19:58:01 -05:00
Dave Horton fae1732ddf update to latest db-helpers 2023-02-07 13:31:48 -05:00
Dave Horton c858624815 Feature/tcp to fs (#72)
* add env var K8S_FEATURE_SERVER_TRANSPORT

* add Contact with private address when using tcp to FS on K8S

* further refinement
2023-01-26 13:33:24 -05:00
13 changed files with 2651 additions and 4233 deletions
+31 -28
View File
@@ -2,16 +2,10 @@ 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:
@@ -20,32 +14,41 @@ jobs:
if: github.event_name == 'push'
steps:
- uses: actions/checkout@v2
- name: Checkout code
uses: actions/checkout@v3
- 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
- name: prepare tag
id: prepare_tag
run: |
IMAGE_ID=ghcr.io/${{ github.repository_owner }}/$IMAGE_NAME
IMAGE_ID=jambonz/sbc-inbound
# 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 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//')
# 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
# Use Docker `latest` tag convention
[ "$VERSION" == "main" ] && VERSION=latest
echo IMAGE_ID=$IMAGE_ID
echo VERSION=$VERSION
echo IMAGE_ID=$IMAGE_ID
echo VERSION=$VERSION
echo "image_id=$IMAGE_ID" >> $GITHUB_OUTPUT
echo "version=$VERSION" >> $GITHUB_OUTPUT
docker tag $IMAGE_NAME $IMAGE_ID:$VERSION
docker push $IMAGE_ID:$VERSION
- name: Login to Docker Hub
uses: docker/login-action@v2
with:
username: ${{ secrets.DOCKERHUB_USERNAME }}
password: ${{ secrets.DOCKERHUB_TOKEN }}
- name: Build and push Docker image
uses: docker/build-push-action@v4
with:
context: .
push: true
tags: ${{ steps.prepare_tag.outputs.image_id }}:${{ steps.prepare_tag.outputs.version }}
build-args: |
GITHUB_REPOSITORY=$GITHUB_REPOSITORY
GITHUB_REF=$GITHUB_REF
+2 -2
View File
@@ -6,8 +6,8 @@ jobs:
build:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v2
- uses: actions/setup-node@v1
- uses: actions/checkout@v3
- uses: actions/setup-node@v3
with:
node-version: 18.x
- run: npm ci
+1 -1
View File
@@ -1,4 +1,4 @@
FROM --platform=linux/amd64 node:18.12.1-alpine3.16 as base
FROM --platform=linux/amd64 node:18.15-alpine3.16 as base
RUN apk --update --no-cache add --virtual .builds-deps build-base python3
+33 -2
View File
@@ -1,10 +1,41 @@
# sbc-inbound ![Build Status](https://github.com/jambonz/sbc-inbound/workflows/CI/badge.svg)
This application provides a part of the SBC (Session Border Controller) functionality of jambonz. It handles incoming INVITE requests from carrier sip trunks or from sip devices and webrtc applications. SIP INVITEs from known carriers are allowed in, while INVITEs from sip devices are challenged to authenticate. SIP traffic that is allowed in is sent on to a jambonz application server in a private subnet.
This application provides a part of the SBC (Session Border Controller) functionality of jambonz platform. It handles incoming INVITE requests from carrier sip trunks or from sip devices and webrtc applications. SIP INVITEs from known carriers are allowed in, while INVITEs from sip devices are challenged to authenticate. SIP traffic that is allowed in is sent on to a jambonz application server in a private subnet.
## Configuration
Configuration is provided via the [npmjs config](https://www.npmjs.com/package/config) package. The following elements make up the configuration for the application:
Configuration is provided via environment variables:
| variable | meaning | required?|
|----------|----------|---------|
|DRACHTIO_HOST| ip address of drachtio server (typically '127.0.0.1')|yes|
|DRACHTIO_PORT| listening port of drachtio server for control connections (typically 9022)|yes|
|DRACHTIO_SECRET| shared secret|yes|
|HTTP_PORT| http listen port |no|
|JAMBONES_LOGLEVEL| log level for application, 'info' or 'debug'|no|
|JAMBONES_MYSQL_HOST| mysql host|yes|
|JAMBONES_MYSQL_PORT| mysql port |no|
|JAMBONES_MYSQL_USER| mysql username|yes|
|JAMBONES_MYSQL_PASSWORD| mysql password|yes|
|JAMBONES_MYSQL_DATABASE| mysql data|yes|
|JAMBONES_MYSQL_CONNECTION_LIMIT| mysql connection limit |no|
|DTMF_LISTEN_PORT| DTMF listening port |no|
|JAMBONES_NG_PROTOCOL| rtpengine NG protocol |no|
|RTPENGINE_PORT| rtpengine port |no|
|JAMBONES_CLUSTER_ID| cluster id |no|
|JAMBONES_NETWORK_CIDR| CIDR of private network that feature server is running in (e.g. '172.31.0.0/16')|yes|
|JAMBONES_REDIS_HOST| redis host|yes|
|JAMBONES_REDIS_PORT|redis port|no|
|JAMBONES_RTPENGINES| commas-separated list of ip:ng-port for rtpengines (e.g. '172.31.32.10:22222')|no|
|JAMBONES_TIME_SERIES_HOST| influxdb host |yes|
|JAMBONES_TIME_SERIES_PORT| influxdb port |no|
|JAMBONES_RECORD_ALL_CALLS| enable auto record calls, 'yes' or 'no' |no|
|K8S| service running as kubernetes service |no|
|K8S_RTPENGINE_SERVICE_NAME| rtpengine service name(required for K8S) |no|
|K8S_FEATURE_SERVER_SERVICE_NAME| feature server service name(required for K8S) |no|
|JWT_SECRET| secret for signing JWT token |yes|
|ENCRYPTION_SECRET| secret for credential encryption(JWT_SECRET is deprecated) |yes|
##### drachtio server location
```
{
+73 -11
View File
@@ -3,10 +3,39 @@ assert.ok(process.env.JAMBONES_MYSQL_HOST &&
process.env.JAMBONES_MYSQL_USER &&
process.env.JAMBONES_MYSQL_PASSWORD &&
process.env.JAMBONES_MYSQL_DATABASE, 'missing JAMBONES_MYSQL_XXX env vars');
assert.ok(process.env.DRACHTIO_PORT || process.env.DRACHTIO_HOST, 'missing DRACHTIO_PORT env var');
if (process.env.JAMBONES_REDIS_SENTINELS) {
assert.ok(process.env.JAMBONES_REDIS_SENTINEL_MASTER_NAME,
'missing JAMBONES_REDIS_SENTINEL_MASTER_NAME env var, JAMBONES_REDIS_SENTINEL_PASSWORD env var is optional');
} else {
assert.ok(process.env.JAMBONES_REDIS_HOST, 'missing JAMBONES_REDIS_HOST env var');
}
assert.ok(process.env.DRACHTIO_PORT || process.env.DRACHTIO_HOST, 'missing DRACHTIO_PORT env vars');
assert.ok(process.env.DRACHTIO_SECRET, 'missing DRACHTIO_SECRET env var');
assert.ok(process.env.JAMBONES_TIME_SERIES_HOST, 'missing JAMBONES_TIME_SERIES_HOST env var');
assert.ok(process.env.JAMBONES_NETWORK_CIDR || process.env.K8S, 'missing JAMBONES_NETWORK_CIDR env var');
const JAMBONES_REDIS_SENTINELS = process.env.JAMBONES_REDIS_SENTINELS ? {
sentinels: process.env.JAMBONES_REDIS_SENTINELS.split(',').map((sentinel) => {
let host, port = 26379;
if (sentinel.includes(':')) {
const arr = sentinel.split(':');
host = arr[0];
port = parseInt(arr[1], 10);
} else {
host = sentinel;
}
return {host, port};
}),
name: process.env.JAMBONES_REDIS_SENTINEL_MASTER_NAME,
...(process.env.JAMBONES_REDIS_SENTINEL_PASSWORD && {
password: process.env.JAMBONES_REDIS_SENTINEL_PASSWORD
}),
...(process.env.JAMBONES_REDIS_SENTINEL_USERNAME && {
username: process.env.JAMBONES_REDIS_SENTINEL_USERNAME
})
} : null;
const Srf = require('drachtio-srf');
const srf = new Srf('sbc-inbound');
const opts = Object.assign({
@@ -28,6 +57,7 @@ const {
commitInterval: 'test' === process.env.NODE_ENV ? 7 : 20
});
const StatsCollector = require('@jambonz/stats-collector');
const CIDRMatcher = require('cidr-matcher');
const stats = new StatsCollector(logger);
const {equalsIgnoreOrder, createHealthCheckApp, systemHealth} = require('./lib/utils');
const {LifeCycleEvents} = require('./lib/constants');
@@ -46,9 +76,11 @@ const {
lookupAccountBySipRealm,
lookupAccountBySid,
lookupAccountCapacitiesBySid,
queryCallLimits
queryCallLimits,
lookupClientByAccountAndUsername
} = require('@jambonz/db-helpers')({
host: process.env.JAMBONES_MYSQL_HOST,
port: process.env.JAMBONES_MYSQL_PORT || 3306,
user: process.env.JAMBONES_MYSQL_USER,
password: process.env.JAMBONES_MYSQL_PASSWORD,
database: process.env.JAMBONES_MYSQL_DATABASE,
@@ -61,8 +93,8 @@ const {
addToSet,
removeFromSet,
incrKey,
decrKey} = require('@jambonz/realtimedb-helpers')({
host: process.env.JAMBONES_REDIS_HOST || 'localhost',
decrKey} = require('@jambonz/realtimedb-helpers')(JAMBONES_REDIS_SENTINELS || {
host: process.env.JAMBONES_REDIS_HOST,
port: process.env.JAMBONES_REDIS_PORT || 6379
}, logger);
@@ -94,7 +126,8 @@ srf.locals = {...srf.locals,
lookupAccountBySid,
lookupAccountBySipRealm,
lookupAccountCapacitiesBySid,
queryCallLimits
queryCallLimits,
lookupClientByAccountAndUsername
},
realtimeDbHelpers: {
createSet,
@@ -107,7 +140,8 @@ const {
getSPForAccount,
wasOriginatedFromCarrier,
getApplicationForDidAndCarrier,
getOutboundGatewayForRefer
getOutboundGatewayForRefer,
getApplicationBySid
} = require('./lib/db-utils')(srf, logger);
srf.locals = {
...srf.locals,
@@ -115,7 +149,8 @@ srf.locals = {
wasOriginatedFromCarrier,
getApplicationForDidAndCarrier,
getOutboundGatewayForRefer,
getFeatureServer: require('./lib/fs-tracking')(srf, logger)
getFeatureServer: require('./lib/fs-tracking')(srf, logger),
getApplicationBySid
};
const activeCallIds = srf.locals.activeCallIds;
@@ -129,7 +164,6 @@ const {
const CallSession = require('./lib/call-session');
if (process.env.DRACHTIO_HOST && !process.env.K8S) {
const CIDRMatcher = require('cidr-matcher');
const cidrs = process.env.JAMBONES_NETWORK_CIDR
.split(',')
.map((s) => s.trim());
@@ -164,6 +198,23 @@ 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});
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.K8S_FEATURE_SERVER_TRANSPORT === 'tcp') {
const matcher = new CIDRMatcher(['192.168.0.0/24', '172.16.0.0/16', '10.0.0.0/8']);
const hostports = hp.split(',');
for (const hp of hostports) {
const arr = /^(.*)\/(.*):(\d+)$/.exec(hp);
if (arr && matcher.contains(arr[2])) {
const hostport = `${arr[2]}:${arr[3]}`;
logger.info(`using sbc private address when sending to feature-server: ${hostport}`);
srf.locals.privateSipAddress = hostport;
}
}
}
});
}
if (process.env.NODE_ENV === 'test') {
srf.on('error', (err) => {
@@ -230,7 +281,8 @@ if (process.env.K8S || process.env.HTTP_PORT) {
if ('test' !== process.env.NODE_ENV) {
/* update call stats periodically */
setInterval(() => {
stats.gauge('sbc.sip.calls.count', activeCallIds.size, ['direction:inbound']);
stats.gauge('sbc.sip.calls.count', activeCallIds.size,
['direction:inbound', `instance_id:${process.env.INSTANCE_ID || 0}`]);
}, 20000);
}
@@ -304,8 +356,18 @@ 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);
logger.info(`got signal ${signal}`);
if (srf.locals.privateSipAddress && setName) {
logger.info(`removing ${srf.locals.privateSipAddress} from set ${setName}`);
removeFromSet(setName, srf.locals.privateSipAddress);
}
if (process.env.K8S) {
lifecycleEmitter.operationalState = LifeCycleEvents.ScaleIn;
if (0 === activeCallIds.size) {
logger.info('exiting immediately since we have no calls in progress');
process.exit(0);
}
}
}
module.exports = {srf, logger};
+3
View File
@@ -57,6 +57,9 @@ module.exports = (logger) => {
}
})();
}
else if (process.env.K8S) {
lifecycleEmitter.scaleIn = () => process.exit(0);
}
return {lifecycleEmitter};
};
+33 -23
View File
@@ -6,15 +6,20 @@ 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 {
SNSClient,
SubscribeCommand,
UnsubscribeCommand } = require('@aws-sdk/client-sns');
const snsClient = new SNSClient({ region: process.env.AWS_REGION, apiVersion: '2010-03-31' });
const {
AutoScalingClient,
DescribeAutoScalingGroupsCommand,
CompleteLifecycleActionCommand } = require('@aws-sdk/client-auto-scaling');
const autoScalingClient = new AutoScalingClient({ region: process.env.AWS_REGION, 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();
@@ -64,7 +69,7 @@ class SnsNotifier extends Emitter {
subscriptionRequestId: this.subscriptionRequestId
}, 'response from SNS SubscribeURL');
const data = await this.describeInstance();
this.lifecycleState = data.AutoScalingInstances[0].LifecycleState;
this.lifecycleState = data.AutoScalingGroups[0].Instances[0].LifecycleState;
this.emit('SubscriptionConfirmation', {publicIp: this.publicIp});
break;
@@ -130,11 +135,12 @@ class SnsNotifier extends Emitter {
async subscribe() {
try {
const response = await sns.subscribe({
const params = {
Protocol: 'http',
TopicArn: process.env.AWS_SNS_TOPIC_ARM,
Endpoint: this.snsEndpoint
}).promise();
};
const response = await snsClient.send(new SubscribeCommand(params));
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}`);
@@ -144,9 +150,10 @@ class SnsNotifier extends Emitter {
async unsubscribe() {
if (!this.subscriptionArn) throw new Error('SnsNotifier#unsubscribe called without an active subscription');
try {
const response = await sns.unsubscribe({
const params = {
SubscriptionArn: this.subscriptionArn
}).promise();
};
const response = await snsClient.send(new UnsubscribeCommand(params));
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}`);
@@ -155,26 +162,29 @@ class SnsNotifier extends Emitter {
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');
});
autoScalingClient.send(new CompleteLifecycleActionCommand(this.scaleInParams))
.then((data) => {
return this.logger.info({data}, 'Successfully completed scale-in action');
})
.catch((err) => {
this.logger.error({err}, 'Error completing scale-in');
});
}
describeInstance() {
return new Promise((resolve, reject) => {
if (!this.instanceId) return reject('instance-id unknown');
autoscaling.describeAutoScalingInstances({
autoScalingClient.send(new DescribeAutoScalingGroupsCommand({
InstanceIds: [this.instanceId]
}, (err, data) => {
if (err) {
}))
.then((data) => {
this.logger.info({data}, 'SnsNotifier: describeInstance');
return resolve(data);
})
.catch((err) => {
this.logger.error({err}, 'Error describing instances');
reject(err);
} else {
this.logger.info({data}, 'SnsNotifier: describeInstance');
resolve(data);
}
});
});
});
}
@@ -188,7 +198,7 @@ module.exports = async function(logger) {
process.on('SIGHUP', async() => {
try {
const data = await notifier.describeInstance();
const state = data.AutoScalingInstances[0].LifecycleState;
const state = data.AutoScalingGroups[0].Instances[0].LifecycleState;
if (state !== notifier.lifecycleState) {
notifier.lifecycleState = state;
switch (state) {
+125 -46
View File
@@ -22,8 +22,10 @@ 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>';
const name = from.name;
const displayName = name ? `${name} ` : '';
if (uri && uri.user) return `${displayName}<sip:${uri.user}@localhost>`;
else return `${displayName}<sip:anonymous@localhost>`;
};
const createSiprecBody = (headers, sdp, type, content) => {
@@ -62,6 +64,7 @@ class CallSession extends Emitter {
this.application_sid = req.locals.application_sid;
this.account_sid = req.locals.account_sid;
this.service_provider_sid = req.locals.service_provider_sid;
this.srsClients = [];
}
get isFromMSTeams() {
@@ -221,7 +224,7 @@ class CallSession extends Emitter {
if (this.req.locals.application_sid) {
Object.assign(headers, {'X-Application-Sid': this.req.locals.application_sid});
}
else if (this.req.authorization) {
if (this.req.authorization) {
if (this.req.authorization.grant && this.req.authorization.grant.application_sid) {
Object.assign(headers, {'X-Application-Sid': this.req.authorization.grant.application_sid});
}
@@ -291,7 +294,6 @@ class CallSession extends Emitter {
} catch (err) {
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);
@@ -316,10 +318,7 @@ class CallSession extends Emitter {
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._stopRecording();
this.srf.endSession(this.req);
});
@@ -328,6 +327,13 @@ class CallSession extends Emitter {
dlg.on('modify', this._onReinvite.bind(this, dlg));
}
_stopRecording() {
if (this.srsClients.length) {
this.srsClients.forEach((c) => c.stop());
this.srsClients = [];
}
}
_setHandlers({uas, uac}) {
this.emit('connected');
const callStart = Date.now();
@@ -379,12 +385,19 @@ class CallSession extends Emitter {
const trunk = ['trunk', 'teams'].includes(this.req.locals.originator) ?
this.req.locals.carrier :
this.req.locals.originator;
const application = await this.srf.locals.getApplicationBySid(application_sid);
const isRecording = this.req.locals.account.record_all_calls || (application && application.record_all_calls);
const day = new Date();
let recording_url = `/Accounts/${this.account_sid}/RecentCalls/${call_sid}/record`;
recording_url += `/${day.getFullYear()}/${(day.getMonth() + 1).toString().padStart(2, '0')}`;
recording_url += `/${day.getDate().toString().padStart(2, '0')}/${this.req.locals.account.record_format}`;
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
trunk,
...(isRecording && {recording_url})
};
this.logger.info({cdr}, 'going to write a cdr now..');
this.writeCdrs({...this.req.locals.cdr,
@@ -392,7 +405,8 @@ class CallSession extends Emitter {
termination_reason: dlg.type === 'uas' ? 'caller hungup' : 'called party hungup',
sip_status: 200,
duration: Math.floor((now - callStart) / 1000),
trunk
trunk,
...(isRecording && {recording_url})
})
.then(() => this.logger.debug('successfully wrote cdr'))
.catch((err) => this.logger.error({err}, 'Error writing cdr for completed call'));
@@ -403,10 +417,7 @@ class CallSession extends Emitter {
dlg.other = null;
other.other = null;
if (this.srsClient) {
this.srsClient.stop();
this.srsClient = null;
}
this._stopRecording();
this.logger.info(`call ended with normal termination, there are ${this.activeCallIds.size} active`);
this.srf.endSession(this.req);
@@ -426,6 +437,12 @@ class CallSession extends Emitter {
// default forwarding of other request types
forwardInDialogRequests(uas, ['notify', 'options', 'message']);
// we need special handling for invite with null sdp followed by 3pcc re-invite
if (uas.local.sdp.includes('a=recvonly') || uas.local.sdp.includes('a=inactive')) {
this.logger.info('incoming call is recvonly or inactive, waiting for re-invite');
this._recvonly = true;
}
}
async _onDTMF(dlg, payload) {
@@ -528,24 +545,70 @@ Duration=${payload.duration} `
}
async _onReinvite(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 reason = req.get('X-Reason');
const isReleasingMedia = reason && dlg.type === 'uac' && ['release-media', 'anchor-media'].includes(reason);
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'];
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});
if (dlg.type === 'uas' && this._recvonly) {
/* seen this from Broadworks - initial INVITE has no SDP, then reINVITE with SDP */
this._recvonly = false; //one-time only
const myMungedSdp = dlg.local.sdp.replace('a=recvonly', 'a=sendrecv').replace('a=inactive', 'a=sendrecv');
this.logger.info({myMungedSdp}, '_onReinvite (3gpp): got a reINVITE with no SDP while in recvonly mode');
res.send(200,
{
body: myMungedSdp
},
(err, req) => {},
async(ack) => {
const remoteOffer = ack.body;
this.logger.info({remoteOffer}, '_onReinvite (3gpp): got ACK for reINVITE with SDP');
let opts = {
...this.rtpEngineOpts.common,
...offerMedia,
'from-tag': fromTag,
'to-tag': toTag,
direction,
sdp: remoteOffer,
};
let response = await this.offer(opts);
if ('ok' !== response.result) {
res.send(488);
throw new Error(`_onReinvite (3gpp): rtpengine failed: offer: ${JSON.stringify(response)}`);
}
this.logger.info({response}, '_onReinvite (3gpp): response from rtpengine for offer');
const fsSdp = await dlg.other.modify(response.sdp);
opts = {
...this.rtpEngineOpts.common,
...answerMedia,
'from-tag': fromTag,
'to-tag': toTag,
sdp: fsSdp
};
response = await this.answer(opts);
if ('ok' !== response.result) {
res.send(488);
throw new Error(`_onReinvite(3gpp): rtpengine failed answer: ${JSON.stringify(response)}`);
}
}
);
}
else {
this.logger.info('got a reINVITE with no SDP; just respond with our current offer');
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 === 'uac' && ['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');
@@ -558,7 +621,8 @@ Duration=${payload.duration} `
direction,
sdp: offeredSdp,
};
if (reason && opts.flags && !opts.flags.includes('reset')) opts.flags.push('reset');
// Dont reset ICE - causes audiocodes webrtrc to fail with "missing ice-ufrag and ice-pwd in re-invite"
// if (reason && opts.flags && !opts.flags.includes('reset')) opts.flags.push('reset');
let response = await this.offer(opts);
if ('ok' !== response.result) {
@@ -627,6 +691,7 @@ Duration=${payload.duration} `
const to = this.req.getParsedHeader('To');
const aorFrom = from.uri;
const aorTo = to.uri;
const headers = contentType === 'application/json' && req.body ? JSON.parse(req.body) : {};
this.logger.info({to, from}, 'startCallRecording request for a call');
const srsUrl = req.get('X-Srs-Url');
@@ -634,7 +699,7 @@ Duration=${payload.duration} `
const callSid = req.get('X-Call-Sid');
const accountSid = req.get('X-Account-Sid');
const applicationSid = req.get('X-Application-Sid');
if (this.srsClient) {
if (this.srsClients.length) {
res.send(400);
this.logger.info('discarding duplicate startCallRecording request for a call');
return;
@@ -644,13 +709,14 @@ Duration=${payload.duration} `
res.send(400);
return;
}
this.srsClient = new SrsClient(this.logger, {
const arr = srsUrl.split(',');
this.srsClients = arr.map((url) => new SrsClient(this.logger, {
srf: dlg.srf,
direction: 'inbound',
originalInvite: this.req,
callingNumber: this.req.callingNumber,
calledNumber: this.req.calledNumber,
srsUrl,
srsUrl: url,
srsRecordingId,
callSid,
accountSid,
@@ -665,42 +731,51 @@ Duration=${payload.duration} `
del: this.del,
blockMedia: this.blockMedia,
unblockMedia: this.unblockMedia,
unsubscribe: this.unsubscribe
});
unsubscribe: this.unsubscribe,
headers
}));
try {
succeeded = await this.srsClient.start();
succeeded = (await Promise.all(
this.srsClients.map((c) => c.start())
)).every((r) => r);
} catch (err) {
this.logger.error({err}, 'Error starting SipRec call recording');
}
}
else if (reason === 'stopCallRecording') {
if (!this.srsClient) {
if (!this.srsClients.length) {
res.send(400);
this.logger.info('discarding stopCallRecording request because we are not recording');
return;
}
try {
succeeded = await this.srsClient.stop();
succeeded = (await Promise.all(
this.srsClients.map((c) => c.stop())
)).every((r) => r);
} catch (err) {
this.logger.error({err}, 'Error stopping SipRec call recording');
}
this.srsClient = null;
this.srsClients = [];
}
else if (reason === 'pauseCallRecording') {
if (!this.srsClient || this.srsClient.paused) {
if (!this.srsClients.length || this.srsClients.every((c) => c.paused)) {
this.logger.info('discarding invalid pauseCallRecording request');
res.send(400);
return;
}
succeeded = await this.srsClient.pause();
succeeded = (await Promise.all(
this.srsClients.map((c) => c.pause())
)).every((r) => r);
}
else if (reason === 'resumeCallRecording') {
if (!this.srsClient || !this.srsClient.paused) {
if (!this.srsClients.length || !this.srsClients.every((c) => c.paused)) {
res.send(400);
this.logger.info('discarding invalid resumeCallRecording request');
return;
}
succeeded = await this.srsClient.resume();
succeeded = (await Promise.all(
this.srsClients.map((c) => c.resume())
)).every((r) => r);
}
res.send(succeeded ? 200 : 503);
}
@@ -753,8 +828,8 @@ Duration=${payload.duration} `
res.send(response.status, {headers: responseHeaders, body: response.body});
}
} catch (err) {
if (this.srsClient) {
this.srsClient = null;
if (this.srsClients.length) {
this.srsClients = [];
}
res.send(500);
this.logger.info({err}, `Error handing INFO request on ${dlg.type} leg`);
@@ -829,6 +904,7 @@ Duration=${payload.duration} `
this.uac = uac;
uac.other = this.uas;
this.uas.other = uac;
uac.on('info', this._onInfo.bind(this, uac));
uac.on('modify', this._onReinvite.bind(this, uac));
uac.on('refer', this._onFeatureServerTransfer.bind(this, uac));
uac.on('destroy', () => {
@@ -838,19 +914,22 @@ Duration=${payload.duration} `
uac.other.destroy();
this.srf.endSession(this.req);
});
// now we can destroy the old dialog
dlg.destroy().catch(() => {});
// modify rtpengine to stream to new feature server
const opts = Object.assign({sdp: uac.remote.sdp, 'to-tag': res.getParsedHeader('To').params.tag},
this.rtpEngineOpts.answer);
const opts = {
...this.rtpEngineOpts.common,
'from-tag': this.rtpEngineOpts.uas.tag,
'to-tag': this.rtpEngineOpts.uac.tag,
sdp: uac.remote.sdp,
flags: ['port latching']
};
const response = await this.answer(opts);
if ('ok' !== response.result) {
res.send(488);
throw new Error(`_onFeatureServerTransfer: rtpengine failed: ${JSON.stringify(response)}`);
throw new Error(`_onFeatureServerTransfer: rtpengine answer failed: ${JSON.stringify(response)}`);
}
dlg.destroy().catch(() => {});
this.logger.info('successfully moved call to new feature server');
} catch (err) {
res.send(488);
this.logger.error(err, 'Error handling refer from feature server');
}
}
+46 -15
View File
@@ -11,6 +11,8 @@ 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.account_sid = acc.account_sid
AND vc.is_active = 1
AND sg.inbound = 1
AND sg.voip_carrier_sid = vc.voip_carrier_sid`;
const sqlSelectAllCarriersForSPByRealm =
@@ -20,6 +22,8 @@ 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 vc.is_active = 1
AND sg.inbound = 1
AND sg.voip_carrier_sid = vc.voip_carrier_sid`;
const sqlSelectAllGatewaysForSP =
@@ -28,7 +32,8 @@ vc.account_sid, vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ip
FROM sip_gateways sg, voip_carriers vc
WHERE sg.voip_carrier_sid = vc.voip_carrier_sid
AND vc.service_provider_sid IS NOT NULL
AND vc.is_active = 1`;
AND vc.is_active = 1
AND sg.inbound = 1`;
const sqlCarriersForAccountBySid =
`SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.account_sid,
@@ -36,10 +41,13 @@ 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.account_sid = ?
AND vc.account_sid = acc.account_sid
AND vc.is_active = 1
AND sg.inbound = 1
AND sg.voip_carrier_sid = vc.voip_carrier_sid`;
const sqlAccountByRealm = 'SELECT * from accounts WHERE sip_realm = ?';
const sqlAccountBySid = 'SELECT * from accounts WHERE account_sid = ?';
const sqlApplicationBySid = 'SELECT * from applications WHERE application_sid = ?';
const sqlQueryApplicationByDid = `
SELECT * FROM phone_numbers
@@ -67,11 +75,15 @@ AND vc.is_active = 1
AND vc.register_sip_realm = ?
AND vc.register_username = ?`;
const gatewayMatchesSourceAddress = (source_address, gw) => {
const gatewayMatchesSourceAddress = (logger, source_address, gw) => {
if (32 === gw.netmask && gw.ipv4 === source_address) return true;
if (gw.netmask < 32) {
const matcher = new CIDRMatcher([`${gw.ipv4}/${gw.netmask}`]);
return matcher.contains(source_address);
try {
const matcher = new CIDRMatcher([`${gw.ipv4}/${gw.netmask}`]);
return matcher.contains(source_address);
} catch (err) {
logger.info({err, gw}, 'gatewayMatchesSourceAddress: Error parsing netmask');
}
}
return false;
};
@@ -80,6 +92,12 @@ module.exports = (srf, logger) => {
const {pool} = srf.locals.dbHelpers;
const pp = pool.promise();
const getApplicationBySid = async(application_sid) => {
const [r] = await pp.query(sqlApplicationBySid, [application_sid]);
if (0 === r.length) return null;
return r[0];
};
const getSPForAccount = async(account_sid) => {
const [r] = await pp.query(sqlSelectSPForAccount, [account_sid]);
if (0 === r.length) return null;
@@ -136,12 +154,12 @@ module.exports = (srf, logger) => {
*/
/* get all the carriers and gateways for the account owning this sip realm */
const [gwAcc] = await pp.query(sqlSelectAllCarriersForAccountByRealm, uri.host);
const [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));
const selected = gw.find(gatewayMatchesSourceAddress.bind(null, logger, req.source_address));
if (selected) {
const [a] = await pp.query(sqlAccountByRealm, uri.host);
const [a] = await pp.query(sqlAccountByRealm, [uri.host]);
if (0 === a.length) return failure;
return {
fromCarrier: true,
@@ -160,14 +178,14 @@ module.exports = (srf, logger) => {
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));
const matches = gw.filter(gatewayMatchesSourceAddress.bind(null, logger, 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);
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 {
@@ -210,7 +228,19 @@ module.exports = (srf, logger) => {
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));
let matches = gw
.filter(gatewayMatchesSourceAddress.bind(null, logger, req.source_address))
.map((gw) => {
return {
voip_carrier_sid: gw.voip_carrier_sid,
name: gw.name,
service_provider_sid: gw.service_provider_sid,
account_sid: gw.account_sid,
application_sid: gw.application_sid
};
});
/* remove duplicates, winnow down to voip_carriers, not gateways */
matches = [...new Set(matches.map(JSON.stringify))].map(JSON.parse);
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(',');
@@ -234,7 +264,7 @@ module.exports = (srf, logger) => {
}
else if (accountLevelGateways.length === 1) {
const [accounts] = await pp.query('SELECT * from accounts where account_sid = ?',
accountLevelGateways[0].account_sid);
[accountLevelGateways[0].account_sid]);
return {
fromCarrier: true,
gateway: accountLevelGateways[0],
@@ -249,11 +279,11 @@ module.exports = (srf, logger) => {
- 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);
[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);
[matches[0].service_provider_sid]);
return {
fromCarrier: true,
gateway: matches[0],
@@ -275,7 +305,7 @@ module.exports = (srf, logger) => {
/* 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);
const [accounts] = await pp.query(sqlAccountBySid, [r[0].account_sid]);
assert(accounts.length);
return {
fromCarrier: true,
@@ -294,6 +324,7 @@ module.exports = (srf, logger) => {
wasOriginatedFromCarrier,
getApplicationForDidAndCarrier,
getOutboundGatewayForRefer,
getSPForAccount
getSPForAccount,
getApplicationBySid
};
};
+3 -1
View File
@@ -15,7 +15,9 @@ module.exports = (srf, logger) => {
return async() => {
try {
if (process.env.K8S) {
return process.env.K8S_FEATURE_SERVER_SERVICE_NAME;
return process.env.K8S_FEATURE_SERVER_TRANSPORT ?
`${process.env.K8S_FEATURE_SERVER_SERVICE_NAME};transport=${process.env.K8S_FEATURE_SERVER_TRANSPORT}` :
process.env.K8S_FEATURE_SERVER_SERVICE_NAME;
}
else {
const fs = await retrieveSet(setName);
+12 -43
View File
@@ -1,8 +1,8 @@
const debug = require('debug')('jambonz:sbc-inbound');
const assert = require('assert');
const Emitter = require('events');
const parseUri = require('drachtio-srf').parseUri;
const {nudgeCallCounts, roundTripTime} = require('./utils');
const digestChallenge = require('@jambonz/digest-utils');
const msProxyIps = process.env.MS_TEAMS_SIP_PROXY_IPS ?
process.env.MS_TEAMS_SIP_PROXY_IPS.split(',').map((i) => i.trim()) :
[];
@@ -22,42 +22,8 @@ const initCdr = (req) => {
};
module.exports = function(srf, logger) {
class AuthOutcomeReporter extends Emitter {
constructor(stats) {
super();
this
.on('regHookOutcome', ({rtt, status}) => {
stats.histogram('app.hook.response_time', rtt, ['hook_type:auth', `status:${status}`]);
})
.on('error', async(err, req) => {
const {account_sid, account} = req.locals;
const {writeAlerts, AlertType} = req.srf.locals;
if (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};
}
else if (err.code === 'ENOTFOUND') {
opts = {...opts, alert_type: AlertType.WEBHOOK_CONNECTION_FAILURE, url: err.hook};
}
else if (err.name === 'StatusError') {
opts = {...opts, alert_type: AlertType.WEBHOOK_STATUS_FAILURE, url: err.hook, status: err.statusCode};
}
if (opts.alert_type) {
try {
await writeAlerts(opts);
} catch (err) {
logger.error({err, opts}, 'Error writing alert');
}
}
}
});
}
}
const {
lookupAuthHook,
lookupAppByTeamsTenant,
lookupAccountBySipRealm,
lookupAccountBySid,
@@ -65,10 +31,6 @@ module.exports = function(srf, logger) {
queryCallLimits
} = srf.locals.dbHelpers;
const {stats, writeCdrs} = srf.locals;
const authenticator = require('@jambonz/http-authenticator')(lookupAuthHook, logger, {
blacklistUnknownRealms: true,
emitter: new AuthOutcomeReporter(stats)
});
const initLocals = (req, res, next) => {
const callId = req.get('Call-ID');
@@ -167,7 +129,7 @@ module.exports = function(srf, logger) {
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');
logger.info({gateway}, 'identifyAccount: incoming call from gateway');
let sid;
if (siprec) {
@@ -194,7 +156,7 @@ module.exports = function(srf, logger) {
};
}
else if (msProxyIps.includes(req.source_address)) {
logger.debug({source_address: req.source_address}, 'identifyAccount: incoming call from Microsoft Teams');
logger.info({source_address: req.source_address}, 'identifyAccount: incoming call from Microsoft Teams');
const uri = parseUri(req.uri);
const app = await lookupAppByTeamsTenant(uri.host);
@@ -216,7 +178,7 @@ module.exports = function(srf, logger) {
else {
req.locals.originator = 'user';
const uri = parseUri(req.uri);
logger.debug({source_address: req.source_address, realm: uri.host},
logger.info({source_address: req.source_address, realm: uri.host},
'identifyAccount: incoming user call');
const account = await lookupAccountBySipRealm(uri.host);
if (!account) {
@@ -240,6 +202,13 @@ module.exports = function(srf, logger) {
account,
application_sid: account.device_calling_application_sid,
webhook_secret: account.webhook_secret,
realm: uri.host,
...(account.registration_hook && {
registration_hook_url: account.registration_hook.url,
registration_hook_method: account.registration_hook.method,
registration_hook_username: account.registration_hook.username,
registration_hook_password: account.registration_hook.password
}),
...req.locals
};
}
@@ -382,7 +351,7 @@ module.exports = function(srf, logger) {
try {
/* TODO: check if this is a gateway that we have an ACL for */
if (req.locals.originator !== 'user') return next();
return authenticator(req, res, next);
return digestChallenge(req, res, next);
} catch (err) {
stats.increment('sbc.terminations', ['sipStatus:500']);
logger.error(err, `${req.get('Call-ID')} Error looking up related info for inbound call`);
+2280 -4053
View File
File diff suppressed because it is too large Load Diff
+9 -8
View File
@@ -1,6 +1,6 @@
{
"name": "sbc-inbound",
"version": "v0.7.8",
"version": "0.8.4",
"main": "app.js",
"engines": {
"node": ">= 12.0.0"
@@ -20,20 +20,21 @@
},
"scripts": {
"start": "node app",
"test": "NODE_ENV=test HTTP_PORT=3050 JAMBONES_NETWORK_CIDR='127.0.0.1/32' JAMBONES_HOSTING=1 SBC_ACCOUNT_SID=ed649e33-e771-403a-8c99-1780eabbc803 JAMBONES_TIME_SERIES_HOST=127.0.0.1 JAMBONES_MYSQL_HOST=127.0.0.1 JAMBONES_MYSQL_USER=jambones_test JAMBONES_MYSQL_PASSWORD=jambones_test JAMBONES_MYSQL_DATABASE=jambones_test JAMBONES_REDIS_HOST=localhost JAMBONES_REDIS_PORT=16379 JAMBONES_LOGLEVEL=error DRACHTIO_SECRET=cymru DRACHTIO_HOST=127.0.0.1 DRACHTIO_PORT=9060 JAMBONES_RTPENGINES=127.0.0.1:12222 JAMBONES_FEATURE_SERVERS=172.38.0.11 node test/ ",
"test": "NODE_ENV=test HTTP_PORT=3050 JAMBONES_NETWORK_CIDR='127.0.0.1/32' JAMBONES_HOSTING=1 JWT_SECRET=foobarbazzle 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.7.3",
"@jambonz/http-authenticator": "^0.2.2",
"@jambonz/db-helpers": "^0.9.1",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/realtimedb-helpers": "^0.6.3",
"@jambonz/realtimedb-helpers": "^0.8.6",
"@jambonz/rtpengine-utils": "^0.4.3",
"@jambonz/siprec-client-utils": "^0.2.0",
"@jambonz/stats-collector": "^0.1.6",
"@jambonz/siprec-client-utils": "^0.2.6",
"@jambonz/stats-collector": "^0.1.9",
"@jambonz/time-series": "^0.2.5",
"aws-sdk": "^2.1261.0",
"@jambonz/digest-utils": "^0.0.3",
"@aws-sdk/client-sns": "^3.360.0",
"@aws-sdk/client-auto-scaling": "^3.360.0",
"bent": "^7.3.12",
"cidr-matcher": "^2.1.1",
"debug": "^4.3.4",