Compare commits

...
86 Commits
Author SHA1 Message Date
Hoan Luu Huu e00379d119 update media session to callee (#78) 2023-04-13 07:43:29 -04:00
EgleH 5174e140a2 Update base image to node:18.15-alpine3.16 (#77) 2023-04-12 13:18:19 -04:00
Dave Horton 7923ececaf push to docker 2023-04-10 09:43:25 -04:00
Anton Voylenko 66974c467d JAMBONES_MYSQL_PORT env added (#76) 2023-04-08 12:16:46 -04:00
Hoan Luu Huu 95bd061e2d feat: update stat collector version (#75) 2023-04-05 12:04:00 -04:00
Anton Voylenko 2f53c0bb4f Update README and add validation (#74)
* update README

* add JAMBONES_TIME_SERIES_HOST validation
2023-04-01 17:51:32 -04:00
Dave Horton f03308b11d bump version 2023-03-28 14:16:20 -04:00
Hoan Luu HuuandQuan HL 7c91ac9d68 feat: add instance_id tag to active call metric (#71)
Co-authored-by: Quan HL <quanluuhoang8@gmail.com>
2023-03-21 07:55:22 -04:00
Dave Horton fcffa1041c refactor of speech-utils 2023-03-14 09:54:46 -04:00
Dave Horton a6d5b25e40 bump version 2023-02-24 10:06:57 -05:00
Snyk bot 5cba07ee44 fix: Dockerfile to reduce vulnerabilities (#69)
The following vulnerabilities are fixed with an upgrade:
- https://snyk.io/vuln/SNYK-UPSTREAM-NODE-3326666
- https://snyk.io/vuln/SNYK-UPSTREAM-NODE-3326668
- https://snyk.io/vuln/SNYK-UPSTREAM-NODE-3326685
- https://snyk.io/vuln/SNYK-UPSTREAM-NODE-3326686
- https://snyk.io/vuln/SNYK-UPSTREAM-NODE-3326688
2023-02-23 07:23:07 -05:00
EgleH 579e21b5bf Upgrade node to node:18.14.0-alpine3.16 (#68) 2023-02-21 07:55:27 -05:00
Dave Horton 8864ab1430 update siprec client with fix for stop recording 2023-02-17 12:11:51 -05:00
Dave Horton 2187b03ca6 update siprec-client with fix for smart tap xml 2023-02-14 14:32:48 -05:00
Dave Horton 0808eeeda4 update to latest siprec client 2023-02-13 16:50:28 -05:00
Dave Horton f841e88bef bump version 2023-02-13 09:15:26 -05:00
Dave Horton b7873654b8 #63: dont select a carrier from a different account in LCR (#64)
* #63: dont select a carrier from a different account in LCR

* handle SIGTERM in K8S

* handle SIGTERM in K8S
2023-02-08 14:30:54 -05:00
Dave Horton b6f9f214c3 fix test case 2023-02-01 09:35:08 -05:00
Dave Horton dbb93732bf bugfix: calls in K8S may have request-uri like sip:<num>:sbc-sip and that's ok 2023-02-01 09:25:12 -05:00
two56andMatt Preskett a553105e55 If running stop srs (#60)
client on destruction

Co-authored-by: Matt Preskett <matt.preskett@netcall.com>
2023-01-23 12:35:41 -05:00
Hoan Luu HuuandQuan HL 36b2c752da fix: B2B uac parse Signal 0 (#59)
Co-authored-by: Quan HL <quanluuhoang8@gmail.com>
2023-01-18 10:33:01 -05:00
Dave Horton f6c0ee6c0d bump version 2023-01-12 16:19:43 -05:00
Dave Horton 96daad8ea1 bugfix: fix ref error on this.req.locals.account.disable_cdrs 2023-01-12 16:04:30 -05:00
Dave Horton b68edf425d minor changes in subscribe dtmf (#57)
* minor changes in subscribe dtmf

* update to rtpengine-utils with fix for not discarding dtmf as dups
2023-01-12 13:17:40 -05:00
Dave Horton 6d44c6c986 gh: run tests on PR 2023-01-09 10:13:53 -05:00
Dave Horton 6d593cbc7d add optional env var to specify udp localPort to bind to 2022-12-30 10:41:39 -05:00
Dave Horton 35b15c1f2c update deps 2022-12-29 09:58:44 -05:00
Dave Horton 7f1d1d61db update to drachtio-srf@4.5.21 with some perf fixes 2022-12-29 09:52:06 -05:00
Dave Horton 20c17fd723 update drachtio-srf 2022-12-28 11:10:15 -06:00
Dave Horton da6fa1da1b bump version 2022-12-24 12:07:10 -06:00
Dave Horton f74dff3b59 bugfix: when releasing media we were restarting ICE which we dont want to do 2022-12-19 21:28:43 -05:00
Dave Horton d9c4e01c36 add env JAMBONES_RECORD_ALL_CALLS to enable global call recording 2022-12-02 13:53:29 -05:00
Dave Horton 3d902c65a4 minor logging 2022-12-01 14:05:23 -05:00
Dave Horton d641504797 fix test 2022-11-29 13:05:51 -05:00
Dave Horton b871812a70 bugfix: was incorrectly making additional attempts after caller hung up or other failures on first attempt 2022-11-28 16:45:21 -05:00
Guilherme RauenandGuilherme Rauen e123a2ef88 update node image to the latest and most secure (#56)
Co-authored-by: Guilherme Rauen <g.rauen@cognigy.com>
2022-11-11 11:26:42 -05:00
Dave Horton 775e63518a include application_sid in cdr 2022-11-07 18:34:13 -05:00
Dave Horton 307cf9bd65 update deps 2022-11-02 13:38:02 -04:00
Dave Horton f48ca4821e update db-helpers 2022-11-01 21:22:55 -04:00
Dave Horton 43ed2bafe1 strip X-Preferred-From-User and X-Preferred-From-Host from outgoing call 2022-10-26 09:35:11 -04:00
Dave Horton 5ee041bdee feature: specify user or host part of From on outdial 2022-10-23 15:27:03 -04:00
Dave Horton 5bbe6a8752 update package-lock.json 2022-10-23 12:20:12 -04:00
Dave Horton 2e5d609bab update rtpengine-utils again 2022-10-23 11:27:34 -04:00
Dave Horton 8f598bf7b0 update to latest rtpengine-utils 2022-10-22 22:26:02 -04:00
Markus FrindtandMarkus Frindt c7d717b3ee [snyk] fix vulnerability (#55)
Co-authored-by: Markus Frindt <m.frindt@cognigy.com>
2022-10-20 21:36:07 -04:00
Dave Horton e9209b37ca include X-Account-Sid header when moving call between FS 2022-10-15 10:59:29 -04:00
Dave Horton 3b98dc6ec2 Bugfix/release media investigation (#54)
* added logging

* when releasing media, include asymettric flag on offer

* update package-lock.json

* remove media handover flag when releasing media
2022-10-14 12:46:14 -04:00
Dave Horton 1c84dd799c update time-series 2022-10-10 09:16:42 +01:00
Dave Horton 4338ae9411 bugfix: inject DMTF flag was inserted over and over 2022-10-07 11:52:36 +01:00
Dave Horton e17e7dbddd add support for connecting to rtpengine via ws 2022-09-27 09:43:05 +01:00
Dave Horton 254479e289 Feature/app call count tracking (#52)
* add call count tracking at the app level (optional)

* update call_counts_app schema

* update rtpengine-utils
2022-09-22 23:44:58 +02:00
Dave Horton a10a311dcb only track service provider calls if JAMBONES_TRACK_SP_CALLS is set 2022-09-20 16:42:22 +02:00
Dave Horton 806cb89c37 include application_sid in cdr 2022-09-20 14:00:33 +02:00
Dave Horton fffa2748d1 Feature/sp limits (#51)
* add account and service provider call limits

* add custom headers when rejecting calls due to max calls limit
2022-09-20 13:13:24 +02:00
Dave Horton 76625c7596 Feature/recent calls enhancement with sp (#50)
* write cdrs with service_provider_sid

* write call counts at the SP level
2022-09-16 13:45:33 +02:00
Paulo Tellesandp.souza 3b0f7ff6eb change node image (#48)
Co-authored-by: p.souza <p.souza@cognigy.com>
2022-09-07 13:20:13 +02:00
Dave Horton 33c75acc9e bump version to 0.7.6 2022-08-26 20:09:43 +02:00
xquanluu bae9ef7638 feat: update time-series 0.11.12 (#47) 2022-08-19 16:21:11 +02:00
Dave Horton 4ed4b38301 update time-series 2022-08-19 09:57:48 +02:00
Dave Horton 6b6f89264f minor logging change 2022-08-17 16:11:10 +02:00
Dave Horton f0e0fba2f1 make rtpengine transcode if non-preferred codec is selected by far end 2022-08-17 14:14:45 +02:00
Dave Horton 290723f234 initial changes for sip info dtmf (#46) 2022-08-11 15:13:07 +02:00
Dave Horton 5a14aa807a update to latest @jambonz/siprec-utils with fix for sdp version 2022-08-08 16:27:59 +02:00
Dave Horton 2505a36db6 fix 2022-08-05 10:11:21 +01:00
Dave Horton 90818206f5 Dockerfile: update base image 2022-08-05 10:10:24 +01:00
Dave Horton 1de4db6ebc Dockerfile: update base image 2022-07-28 12:58:23 +01:00
Dave Horton 7d2125788f when releasing media, use asymetric flag so that rtpengine does react to a spurious final packet from freeswitch by incorrectly sending rtp there 2022-07-26 12:25:58 +01:00
Snyk bot 3b83c1bda8 fix: upgrade drachtio-srf from 4.5.0 to 4.5.1 (#42)
Snyk has created this PR to upgrade drachtio-srf from 4.5.0 to 4.5.1.

See this package in npm:
https://www.npmjs.com/package/drachtio-srf

See this project in Snyk:
https://app.snyk.io/org/davehorton/project/282b1881-3e13-4fad-85bd-3b1c662c2134?utm_source=github&utm_medium=referral&page=upgrade-pr
2022-07-13 09:16:37 +02:00
Paulo Tellesandp.souza f352bf885c improve dockerfile to fix snyk security issues (#41)
Co-authored-by: p.souza <p.souza@cognigy.com>
2022-07-07 15:19:14 +02:00
Dave Horton 08a4f5defb update Dockerfile 2022-06-23 09:57:38 -04:00
Dave Horton 87df38110b update to azure 1.22.0 2022-06-11 16:23:08 -04:00
Dave Horton 06e370fa59 update deps 2022-06-11 11:55:21 -04:00
Dave Horton e4ed2cea26 healthcheck improvements (#38) 2022-04-12 15:45:33 -04:00
Dave Horton 0c8967acdb bump version 2022-04-06 08:18:22 -04:00
Dave Horton 1dfeed3ac8 bugfix: mem leak 2022-04-05 19:57:03 -04:00
Dave Horton d10bda2926 update deps 2022-04-01 13:46:16 -04:00
Dave Horton 3938773738 mem leak fix 2022-03-28 09:09:14 -04:00
Dave Horton f881002943 write otel trace_id to call history 2022-03-23 09:26:04 -04:00
Dave Horton 2266b80e73 release drachtio connection on final outbound failure 2022-03-18 14:34:30 -04:00
Snyk bot 4c3d6ddf0c fix: Dockerfile to reduce vulnerabilities (#35)
The following vulnerabilities are fixed with an upgrade:
- https://snyk.io/vuln/SNYK-DEBIAN11-GNUTLS28-2419151
- https://snyk.io/vuln/SNYK-DEBIAN11-OPENSSL-2388380
- https://snyk.io/vuln/SNYK-DEBIAN11-OPENSSL-2426309
- https://snyk.io/vuln/SNYK-DEBIAN11-UTILLINUX-2401081
- https://snyk.io/vuln/SNYK-DEBIAN11-UTILLINUX-2401081
2022-03-18 07:55:10 -04:00
Dave Horton 3d1bcb23f4 bump version 2022-03-08 20:16:19 -05:00
Dave Horton 98ecfa20aa add support for redis auth 2022-03-08 09:30:44 -05:00
Dave Horton 56efe50aec forward incoming REFER to FS (#34) 2022-03-05 15:21:49 -05:00
Dave Horton 891d0ff38b use registered contact as uri when sending to user (#32) 2022-02-23 15:09:18 -05:00
Dave Horton 78d6cb5f22 use registered contact as uri when sending to user (#31) 2022-02-17 21:35:19 -05:00
Dave Horton 23255a71db added pre-commit hook for linting 2022-02-14 13:18:46 -05:00
15 changed files with 3399 additions and 4435 deletions
+1 -1
View File
@@ -8,7 +8,7 @@
"jsx": false,
"modules": false
},
"ecmaVersion": 2018
"ecmaVersion": 2020
},
"plugins": ["promise"],
"rules": {
+4 -6
View File
@@ -1,17 +1,15 @@
name: CI
on:
push:
workflow_dispatch:
on: [push, pull_request]
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: 14.x
node-version: lts/*
- run: npm ci
- run: npm run jslint
- run: npm test
+31 -30
View File
@@ -2,16 +2,8 @@ name: Docker
on:
push:
# Publish `main` as Docker `latest` image.
branches:
- main
# Publish `v1.2.3` tags as releases.
tags:
- v*
env:
IMAGE_NAME: sbc-outbound
- '*'
jobs:
push:
@@ -20,32 +12,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=$GITHUB_REPOSITORY
# 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
+4
View File
@@ -0,0 +1,4 @@
#!/bin/sh
. "$(dirname "$0")/_/husky.sh"
npm run jslint
+19 -6
View File
@@ -1,10 +1,23 @@
FROM node:17.4-slim
FROM --platform=linux/amd64 node:18.15-alpine3.16 as base
RUN apk --update --no-cache add --virtual .builds-deps build-base python3
WORKDIR /opt/app/
COPY package.json ./
RUN npm install
RUN npm prune
COPY . /opt/app
FROM base as build
COPY package.json package-lock.json ./
RUN npm ci
COPY . .
FROM base
COPY --from=build /opt/app /opt/app/
ARG NODE_ENV
ENV NODE_ENV $NODE_ENV
CMD [ "npm", "start" ]
CMD [ "node", "app.js" ]
+14 -14
View File
@@ -1,6 +1,6 @@
# sbc-outbound ![Build Status](https://github.com/jambonz/sbc-outbound/workflows/CI/badge.svg)
This application provides a part of the SBC (Session Border Controller) functionality of jambonz. It handles outbound INVITE requests from the cpaas application server that is going to carrier sip trunks or registered sip users/devices, including webrtc applications.
This application provides a part of the SBC (Session Border Controller) functionality of jambonz platfrom. It handles outbound INVITE requests from the cpaas application server that is going to carrier sip trunks or registered sip users/devices, including webrtc applications.
## Configuration
@@ -11,22 +11,26 @@ Configuration is provided via environment variables:
|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|
|ENABLE_METRICS| if 1, metrics will be generated|no|
|HTTP_PORT| tcp port 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|yes|
|JAMBONES_REDIS_PORT|redis port|no|
|JAMBONES_RTPENGINES| commans-separated list of ip:ng-port for rtpengines (e.g. '172.31.32.10:22222')|yes|
|JAMBONES_SBCS| list of IP addresses (on the internal network) of SBCs, comma-separated|yes|
|STATS_HOST| ip address of metrics host (usually '127.0.0.1' since telegraf is installed locally|no|
|STATS_PORT| listening port for metrics host|no|
|STATS_PROTOCOL| 'tcp' or 'udp'|no|
|STATS_TELEGRAF| if 1, metrics will be generated in telegraf format|no|
|JAMBONES_TIME_SERIES_HOST| influxdb host |yes|
|JAMBONES_RECORD_ALL_CALLS| enable auto record calls |no|
|K8S| service running as kubernetes service |no|
|K8S_RTPENGINE_SERVICE_NAME| rtpengine service name(required for K8S) |no|
### running under pm2
Typically, this application runs under [pm2](https://pm2.io) using an [ecosystem.config.js](https://pm2.keymetrics.io/docs/usage/application-declaration/) file similar to this:
@@ -58,17 +62,13 @@ module.exports = {
JAMBONES_MYSQL_CONNECTION_LIMIT: 10,
JAMBONES_REDIS_HOST: 'jambonz.zzzzzzz.0001.usw1.cache.amazonaws.com',
JAMBONES_REDIS_PORT: 6379,
ENABLE_METRICS: 1,
STATS_HOST: '127.0.0.1',
STATS_PORT: 8125,
STATS_PROTOCOL: 'tcp',
STATS_TELEGRAF: 1,
JAMBONES_TIME_SERIES_HOST: '172.31.32.11',
JAMBONES_NETWORK_CIDR: '172.31.0.0/16'
}
}]
};
```
#### Running the test suite
To run the included test suite, you will need to have docker installed on your laptop.
```
+68 -15
View File
@@ -7,16 +7,20 @@ 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 var');
assert.ok(process.env.DRACHTIO_SECRET, 'missing DRACHTIO_SECRET env var');
assert.ok(process.env.JAMBONES_NETWORK_CIDR || process.env.K8S, 'missing JAMBONES_NETWORK_CIDR env var');
assert.ok(process.env.JAMBONES_TIME_SERIES_HOST, 'missing JAMBONES_TIME_SERIES_HOST env var');
const Srf = require('drachtio-srf');
const srf = new Srf('sbc-outbound');
const CIDRMatcher = require('cidr-matcher');
const {pingMsTeamsGateways, equalsIgnoreOrder} = require('./lib/utils');
const {equalsIgnoreOrder, pingMsTeamsGateways, createHealthCheckApp, systemHealth} = require('./lib/utils');
const opts = Object.assign({
timestamp: () => {return `, "time": "${new Date().toISOString()}"`;}
}, {level: process.env.JAMBONES_LOGLEVEL || 'info'});
const logger = require('pino')(opts);
const {
writeCallCount,
writeCallCountSP,
writeCallCountApp,
writeCdrs,
queryCdrs,
writeAlerts,
@@ -32,21 +36,25 @@ const CallSession = require('./lib/call-session');
const setNameRtp = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-rtp`;
const rtpServers = [];
const {
ping,
performLcr,
lookupAllTeamsFQDNs,
lookupAccountBySipRealm,
lookupAccountBySid,
lookupAccountCapacitiesBySid,
lookupSipGatewaysByCarrier,
lookupCarrierBySid
lookupCarrierBySid,
queryCallLimits
} = 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,
connectionLimit: process.env.JAMBONES_MYSQL_CONNECTION_LIMIT || 10
}, logger);
const {
client: redisClient,
createHash,
retrieveHash,
incrKey,
@@ -54,27 +62,35 @@ const {
retrieveSet,
isMemberOfSet
} = require('@jambonz/realtimedb-helpers')({
host: process.env.JAMBONES_REDIS_HOST || 'localhost',
host: process.env.JAMBONES_REDIS_HOST,
port: process.env.JAMBONES_REDIS_PORT || 6379
}, logger);
const activeCallIds = new Map();
const Emitter = require('events');
const idleEmitter = new Emitter();
srf.locals = {...srf.locals,
stats,
writeCallCount,
writeCallCountSP,
writeCallCountApp,
writeCdrs,
writeAlerts,
AlertType,
queryCdrs,
activeCallIds,
idleEmitter,
dbHelpers: {
ping,
performLcr,
lookupAllTeamsFQDNs,
lookupAccountBySipRealm,
lookupAccountBySid,
lookupAccountCapacitiesBySid,
lookupSipGatewaysByCarrier,
lookupCarrierBySid
lookupCarrierBySid,
queryCallLimits
},
realtimeDbHelpers: {
createHash,
@@ -88,10 +104,12 @@ const {initLocals, checkLimits, route} = require('./lib/middleware')(srf, logger
host: process.env.JAMBONES_REDIS_HOST,
port: process.env.JAMBONES_REDIS_PORT || 6379
});
const ngProtocol = process.env.JAMBONES_NG_PROTOCOL || 'udp';
const ngPort = process.env.RTPENGINE_PORT || ('udp' === ngProtocol ? 22222 : 8080);
const {getRtpEngine, setRtpEngines} = require('@jambonz/rtpengine-utils')([], logger, {
emitter: stats,
//emitter: stats,
dtmfListenPort: process.env.DTMF_LISTEN_PORT || 22225,
protocol: 'udp'
protocol: ngProtocol
});
srf.locals.getRtpEngine = getRtpEngine;
@@ -137,22 +155,41 @@ srf.invite((req, res) => {
session.connect();
});
if (process.env.K8S) {
if (process.env.K8S || process.env.HTTP_PORT) {
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});
}
const getCount = () => srf.locals.activeCallIds.size;
createHealthCheckApp(PORT, logger)
.then((app) => {
healthCheck({
app,
logger,
path: '/',
fn: getCount
});
healthCheck({
app,
logger,
path: '/system-health',
fn: systemHealth.bind(null, redisClient, ping, getCount)
});
return;
})
.catch((err) => {
logger.error({err}, 'Error creating health check server');
});
}
if ('test' !== process.env.NODE_ENV) {
/* update call stats periodically */
setInterval(() => {
stats.gauge('sbc.sip.calls.count', activeCallIds.size, ['direction:outbound']);
}, 5000);
stats.gauge('sbc.sip.calls.count', activeCallIds.size, ['direction:outbound',
`instance_id:${process.env.INSTANCE_ID || 0}`]);
}, 20000);
}
const lookupRtpServiceEndpoints = (lookup, serviceName) => {
logger.debug(`dns lookup for ${serviceName}..`);
lookup(serviceName, {family: 4, all: true}, (err, addresses) => {
if (err) {
logger.error({err}, `Error looking up ${serviceName}`);
@@ -164,7 +201,7 @@ const lookupRtpServiceEndpoints = (lookup, serviceName) => {
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}`));
setRtpEngines(rtpServers.map((a) => `${a}:${ngPort}`));
}
});
};
@@ -191,7 +228,7 @@ else {
logger.debug({newArray, rtpServers}, 'getActiveRtpServers');
if (!equalsIgnoreOrder(newArray, rtpServers)) {
logger.info({newArray}, 'resetting active rtpengines');
setRtpEngines(newArray.map((a) => `${a}:${process.env.RTPENGINE_PORT || 22222}`));
setRtpEngines(newArray.map((a) => `${a}:${ngPort}`));
rtpServers.length = 0;
Array.prototype.push.apply(rtpServers, newArray);
}
@@ -207,4 +244,20 @@ else {
pingMsTeamsGateways(logger, srf);
process.on('SIGUSR2', handle.bind(null));
process.on('SIGTERM', handle.bind(null));
function handle(signal) {
logger.info(`got signal ${signal}`);
if (process.env.K8S) {
if (0 === activeCallIds.size) {
logger.info('exiting immediately since we have no calls in progress');
process.exit(0);
}
else {
idleEmitter.once('idle', () => process.exit(0));
}
}
}
module.exports = {srf};
+276 -40
View File
@@ -1,5 +1,7 @@
const Emitter = require('events');
const {makeRtpEngineOpts, makeCallCountKey} = require('./utils');
const sdpTransform = require('sdp-transform');
const SrsClient = require('@jambonz/siprec-client-utils');
const {makeRtpEngineOpts, nudgeCallCounts} = require('./utils');
const {forwardInDialogRequests} = require('drachtio-fn-b2b-sugar');
const {SipError, stringifyUri, parseUri} = require('drachtio-srf');
const debug = require('debug')('jambonz:sbc-outbound');
@@ -11,10 +13,17 @@ const makeInviteInProgressKey = (callid) => `sbc-out-iip${callid}`;
*/
const createBLegFromHeader = (req, teams) => {
const from = req.getParsedHeader('From');
const host = teams ? req.get('X-MS-Teams-Tenant-FQDN') : 'localhost';
const uri = parseUri(from.uri);
if (uri && uri.user) return `sip:${uri.user}@${host}`;
return `sip:anonymous@${host}`;
let user = uri.user || 'anonymous';
let host = 'localhost';
if (teams) {
host = req.get('X-MS-Teams-Tenant-FQDN');
}
else if (req.has('X-Preferred-From-User') || req.has('X-Preferred-From-Host')) {
user = req.get('X-Preferred-From-User') || user;
host = req.get('X-Preferred-From-Host') || host;
}
return `sip:${user}@${host}`;
};
const createBLegToHeader = (req, teams) => {
const to = req.getParsedHeader('To');
@@ -31,11 +40,13 @@ const initCdr = (srf, req) => {
const to = arr ? arr[1] : req.calledNumber;
arr = regex.exec(req.callingNumber);
const from = arr ? arr[1] : req.callingNumber;
const applicationSid = req.get('X-Application-Sid');
return {
account_sid: req.get('X-Account-Sid'),
call_sid: req.get('X-Call-Sid'),
sip_callid: req.get('Call-ID'),
...(applicationSid && {application_sid: applicationSid}),
from,
to,
duration: 0,
@@ -43,10 +54,20 @@ const initCdr = (srf, req) => {
attempted_at: Date.now(),
direction: 'outbound',
host: srf.locals.sipAddress,
remote_host: uri.host
remote_host: uri.host,
trace_id: req.get('X-Trace-ID') || '00000000000000000000000000000000'
};
};
const updateRtpEngineFlags = (sdp, opts) => {
try {
const parsed = sdpTransform.parse(sdp);
const codec = parsed.media[0].rtp[0].codec;
if (['PCMU', 'PCMA'].includes(codec)) opts.flags.push(`codec-accept-${codec}`);
} catch (err) {}
return opts;
};
class CallSession extends Emitter {
constructor(logger, req, res) {
super();
@@ -56,27 +77,55 @@ class CallSession extends Emitter {
this.logger = logger.child({callId: req.get('Call-ID')});
this.useWss = req.locals.registration && req.locals.registration.protocol === 'wss';
this.stats = this.srf.locals.stats;
this.idleEmitter = this.srf.locals.idleEmitter;
this.activeCallIds = this.srf.locals.activeCallIds;
this.writeCdrs = this.srf.locals.writeCdrs;
this.incrKey = req.srf.locals.realtimeDbHelpers.incrKey;
this.decrKey = req.srf.locals.realtimeDbHelpers.decrKey;
this.callCountKey = makeCallCountKey(req.locals.account_sid);
const {performLcr, lookupCarrierBySid, lookupSipGatewaysByCarrier} = this.srf.locals.dbHelpers;
this.performLcr = performLcr;
this.lookupCarrierBySid = lookupCarrierBySid;
this.lookupSipGatewaysByCarrier = lookupSipGatewaysByCarrier;
this._mediaReleased = false;
}
get account_sid() {
return this.req.locals.account_sid;
}
get application_sid() {
return this.req.locals.application_sid;
}
get privateSipAddress() {
return this.srf.locals.privateSipAddress;
}
get isMediaReleased() {
return this._mediaReleased;
}
get calleeIsUsingSrtp() {
const tp = this.rtpEngineOpts?.uac?.mediaOpts['transport-protocol'];
return tp && -1 !== tp.indexOf('SAVP');
}
subscribeForDTMF(dlg) {
if (!this._subscribedForDTMF) {
this._subscribedForDTMF = true;
this.subscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uac.tag,
this._onDTMF.bind(this, dlg));
}
}
unsubscribeForDTMF() {
if (this._subscribedForDTMF) {
this._subscribedForDTMF = false;
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uac.tag);
}
}
async connect() {
const teams = this.teams = this.req.locals.target === 'teams';
const engine = this.srf.locals.getRtpEngine();
@@ -93,8 +142,12 @@ class CallSession extends Emitter {
unblockMedia,
blockDTMF,
unblockDTMF,
playDTMF,
subscribeDTMF,
unsubscribeDTMF
unsubscribeDTMF,
subscribeRequest,
subscribeAnswer,
unsubscribe
} = engine;
const {createHash, retrieveHash} = this.srf.locals.realtimeDbHelpers;
this.offer = offer;
@@ -104,8 +157,12 @@ class CallSession extends Emitter {
this.unblockMedia = unblockMedia;
this.blockDTMF = blockDTMF;
this.unblockDTMF = unblockDTMF;
this.playDTMF = playDTMF;
this.subscribeDTMF = subscribeDTMF;
this.unsubscribeDTMF = unsubscribeDTMF;
this.subscribeRequest = subscribeRequest;
this.subscribeAnswer = subscribeAnswer;
this.unsubscribe = unsubscribe;
this.rtpEngineOpts = makeRtpEngineOpts(this.req, false, this.useWss || teams, teams);
this.rtpEngineResource = {destroy: this.del.bind(null, this.rtpEngineOpts.common)};
@@ -126,7 +183,7 @@ class CallSession extends Emitter {
if (this.req.locals.registration) {
debug(`sending call to registered user ${JSON.stringify(this.req.locals.registration)}`);
const contact = this.req.locals.registration.contact;
let destUri = this.req.uri;
let destUri = contact;
if (this.req.has('X-Override-To')) {
const dest = this.req.get('X-Override-To');
const uri = parseUri(contact);
@@ -134,12 +191,9 @@ class CallSession extends Emitter {
destUri = stringifyUri(uri);
this.logger.info(`overriding destination user with ${dest}, so final uri is ${destUri}`);
}
if (contact.includes('transport=ws')) {
uris = [contact];
}
else {
uris = [destUri];
if (!contact.includes('transport=ws')) {
proxy = this.req.locals.registration.proxy;
uris = [destUri];
}
}
else if (this.req.locals.target === 'forward') {
@@ -223,13 +277,13 @@ class CallSession extends Emitter {
}
// rtpengine 'offer'
const opts = {
const opts = updateRtpEngineFlags(this.req.body, {
...this.rtpEngineOpts.common,
...this.rtpEngineOpts.uac.mediaOpts,
'from-tag': this.rtpEngineOpts.uas.tag,
direction: ['private', 'public'],
sdp: this.req.body
};
});
const response = await this.offer(opts);
debug(`response from rtpengine to offer ${JSON.stringify(response)}`);
this.logger.debug({offer: opts, response}, 'initial offer to rtpengine');
@@ -289,13 +343,15 @@ class CallSession extends Emitter {
'all',
'-X-MS-Teams-FQDN',
'-X-MS-Teams-Tenant-FQDN',
'X-CID',
'-X-Trace-ID',
'-Allow',
'-Session-Expires',
'-X-Requested-Carrier-Sid',
'-X-Jambonz-Routing',
'-X-Jambonz-FS-UUID',
'Min-SE'
'-X-Preferred-From-User',
'X-Preferred-From-Host',
'-X-Jambonz-FS-UUID',
],
proxyResponseHeaders: [
'all',
@@ -334,10 +390,15 @@ class CallSession extends Emitter {
else if (this.req.locals.registration) trunk = 'user';
else trunk = 'sipUri';
}
if (!this.req.locals.account.disable_cdrs) {
if (this.req.locals.account?.disable_cdrs) {
this.logger.debug('cdrs disabled for this account');
}
else {
this.req.locals.cdr = {
...initCdr(this.req.srf, inv),
service_provider_sid: this.req.locals.service_provider_sid,
account_sid: this.req.locals.account_sid,
...(this.req.locals.application_sid && {application_sid: this.req.locals.application_sid}),
trunk
};
}
@@ -382,6 +443,7 @@ class CallSession extends Emitter {
this.emit('failed');
this.rtpEngineResource.destroy()
.catch((err) => this.logger.info({err}, 'Error destroying rtpe after failure'));
this.srf.endSession(this.req);
const tags = ['accepted:no', `sipStatus:${status}`];
this.stats.increment('sbc.originations', tags);
@@ -389,9 +451,10 @@ class CallSession extends Emitter {
this.writeCdrs({...this.req.locals.cdr,
terminated_at: Date.now(),
termination_reason: 487 === status ? 'caller abandoned' : 'failed',
sip_status: status,
sip_status: status
}).catch((err) => this.logger.error({err}, 'Error writing cdr for call failure'));
}
return;
}
else {
this.logger.info(`got ${err.status}, cranking back to next destination`);
@@ -421,20 +484,29 @@ class CallSession extends Emitter {
this.uas = uas;
this.uac = uac;
[uas, uac].forEach((dlg) => {
dlg.on('destroy', () => {
this.logger.info('call ended');
dlg.on('destroy', async() => {
const other = dlg.other;
this.rtpEngineResource.destroy();
this.activeCallIds.delete(this.req.get('Call-ID'));
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uac.tag);
dlg.other.destroy();
this.unsubscribeForDTMF();
//this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uac.tag);
try {
await other.destroy();
} catch (err) {}
this.decrKey(this.callCountKey)
.then((count) => {
this.logger.debug(`after hangup there are ${count} active calls for this account`);
debug(`after hangup there are ${count} active calls for this account`);
return;
})
.catch((err) => this.logger.error({err}, 'Error decrementing call count'));
const trackingOn = process.env.JAMBONES_TRACK_ACCOUNT_CALLS ||
process.env.JAMBONES_TRACK_SP_CALLS ||
process.env.JAMBONES_TRACK_APP_CALLS;
if (process.env.JAMBONES_HOSTING || trackingOn) {
const {writeCallCount, writeCallCountSP, writeCallCountApp} = this.req.srf.locals;
await nudgeCallCounts(this.logger, {
service_provider_sid: this.service_provider_sid,
account_sid: this.account_sid,
application_sid: this.application_sid
}, this.decrKey, {writeCallCountSP, writeCallCount, writeCallCountApp})
.catch((err) => this.logger.error(err, 'Error decrementing call counts'));
}
/* write cdr for connected call */
if (this.req.locals.cdr) {
@@ -447,17 +519,32 @@ class CallSession extends Emitter {
duration: Math.floor((now - callStart) / 1000)
}).catch((err) => this.logger.error({err}, 'Error writing cdr for completed call'));
}
/* de-link the 2 Dialogs for GC */
dlg.removeAllListeners();
other.removeAllListeners();
dlg.other = null;
other.other = null;
if (this.srsClient) {
this.srsClient.stop();
this.srsClient = null;
}
this.logger.info(`call ended with normal termination, there are ${this.activeCallIds.size} active`);
if (this.activeCallIds.size === 0) this.idleEmitter.emit('idle');
this.srf.endSession(this.req);
});
});
this.subscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uac.tag,
this._onDTMF.bind(this, uas));
this.subscribeForDTMF(uas);
//this.subscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uac.tag,
// this._onDTMF.bind(this, uas));
uas.on('modify', this._onReinvite.bind(this, uas));
uac.on('modify', this._onReinvite.bind(this, uac));
uas.on('refer', this._onFeatureServerTransfer.bind(this, uas));
uac.on('refer', this._onRefer.bind(this, uac));
uas.on('info', this._onInfo.bind(this, uas));
uac.on('info', this._onInfo.bind(this, uac));
@@ -466,6 +553,23 @@ class CallSession extends Emitter {
forwardInDialogRequests(uac, ['notify', 'options', 'message']);
}
async _onRefer(dlg, req, res) {
/* REFER coming in from a sip device, forward to feature server */
try {
const response = await dlg.other.request({
method: 'REFER',
headers: {
'Refer-To': req.get('Refer-To'),
'Referred-By': req.get('Referred-By'),
'User-Agent': req.get('User-Agent')
}
});
res.send(response.status, response.reason);
} catch (err) {
this.logger.error({err}, 'CallSession:_onRefer: error handling incoming REFER');
}
}
async _onDTMF(dlg, payload) {
this.logger.info({payload}, '_onDTMF');
try {
@@ -497,11 +601,17 @@ Duration=${payload.duration} `
async _onReinvite(dlg, req, res) {
try {
const reason = req.get('X-Reason');
const isReleasingMedia = reason && dlg.type === 'uas' && ['release-media', 'anchor-media'].includes(reason);
const fromTag = dlg.type === 'uas' ? this.rtpEngineOpts.uas.tag : this.rtpEngineOpts.uac.tag;
const toTag = dlg.type === 'uas' ? this.rtpEngineOpts.uac.tag : this.rtpEngineOpts.uas.tag;
const offerMedia = dlg.type === 'uas' ? this.rtpEngineOpts.uac.mediaOpts : this.rtpEngineOpts.uas.mediaOpts;
const answerMedia = dlg.type === 'uas' ? this.rtpEngineOpts.uas.mediaOpts : this.rtpEngineOpts.uac.mediaOpts;
const direction = dlg.type === 'uas' ? ['private', 'public'] : ['public', 'private'];
if (isReleasingMedia) {
if (!offerMedia.flags.includes('port latching')) offerMedia.flags.push('port latching');
if (!offerMedia.flags.includes('asymmetric')) offerMedia.flags.push('asymmetric');
offerMedia.flags = offerMedia.flags.filter((f) => f !== 'media handover');
}
let opts = {
...this.rtpEngineOpts.common,
...offerMedia,
@@ -510,23 +620,29 @@ Duration=${payload.duration} `
direction,
sdp: req.body,
};
if (reason && opts.flags && !opts.flags.includes('reset')) opts.flags.push('reset');
// DH: this was restarting ICE, which we don't want to do
//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)}`);
}
this.logger.debug({opts, response}, 'CallSession:_onReinvite: (offer)');
/* if this is a re-invite from the FS to change media anchoring, avoid sending the reinvite out */
let sdp;
if (reason && dlg.type === 'uas' && ['release-media', 'anchor-media'].includes(reason)) {
if (isReleasingMedia && !this.calleeIsUsingSrtp) {
this.logger.info(`got a reinvite from FS to ${reason}`);
sdp = dlg.other.remote.sdp;
if (!answerMedia.flags.includes('port latching')) answerMedia.flags.push('port latching');
if (!answerMedia.flags.includes('asymmetric')) answerMedia.flags.push('asymmetric');
answerMedia.flags = answerMedia.flags.filter((f) => f !== 'media handover');
this._mediaReleased = 'release-media' === reason;
}
else {
sdp = await dlg.other.modify(response.sdp);
this.logger.info({sdp}, 'CallSession:_onReinvite: got sdp from 200 OK to invite we sent');
}
opts = {
...this.rtpEngineOpts.common,
@@ -540,6 +656,7 @@ Duration=${payload.duration} `
res.send(488);
throw new Error(`_onReinvite: rtpengine failed: ${JSON.stringify(response)}`);
}
this.logger.debug({opts, sdp: response.sdp}, 'CallSession:_onReinvite: (answer) sending back upstream');
res.send(200, {body: response.sdp});
} catch (err) {
this.logger.error(err, 'Error handling reinvite');
@@ -548,6 +665,7 @@ Duration=${payload.duration} `
async _onInfo(dlg, req, res) {
try {
const contentType = req.get('Content-Type');
if (dlg.type === 'uas' && req.has('X-Reason')) {
const toTag = this.rtpEngineOpts.uac.tag;
const reason = req.get('X-Reason');
@@ -566,6 +684,123 @@ Duration=${payload.duration} `
const response = Promise.all([this.unblockMedia(opts), this.unblockDTMF(opts)]);
this.logger.info({response}, `_onInfo: response to rtpengine command for ${reason}`);
}
else if (reason.includes('CallRecording')) {
let succeeded = false;
if (reason === 'startCallRecording') {
const from = this.req.getParsedHeader('From');
const to = this.req.getParsedHeader('To');
const aorFrom = from.uri;
const aorTo = to.uri;
this.logger.info({to, from}, 'startCallRecording request for an outbound call');
const srsUrl = req.get('X-Srs-Url');
const srsRecordingId = req.get('X-Srs-Recording-ID');
const callSid = req.get('X-Call-Sid');
const accountSid = req.get('X-Account-Sid');
const applicationSid = req.get('X-Application-Sid');
if (this.srsClient) {
res.send(400);
this.logger.info('discarding duplicate startCallRecording request for a call');
return;
}
if (!srsUrl) {
this.logger.info('startCallRecording request is missing X-Srs-Url header');
res.send(400);
return;
}
this.srsClient = new SrsClient(this.logger, {
srf: dlg.srf,
direction: 'outbound',
originalInvite: this.req,
callingNumber: this.req.callingNumber,
calledNumber: this.req.calledNumber,
srsUrl,
srsRecordingId,
callSid,
accountSid,
applicationSid,
rtpEngineOpts: this.rtpEngineOpts,
toTag,
aorFrom,
aorTo,
subscribeRequest: this.subscribeRequest,
subscribeAnswer: this.subscribeAnswer,
del: this.del,
blockMedia: this.blockMedia,
unblockMedia: this.unblockMedia,
unsubscribe: this.unsubscribe
});
try {
succeeded = await this.srsClient.start();
} catch (err) {
this.logger.error({err}, 'Error starting SipRec call recording');
}
}
else if (reason === 'stopCallRecording') {
if (!this.srsClient) {
res.send(400);
this.logger.info('discarding stopCallRecording request because we are not recording');
return;
}
try {
succeeded = await this.srsClient.stop();
} catch (err) {
this.logger.error({err}, 'Error stopping SipRec call recording');
}
this.srsClient = null;
}
else if (reason === 'pauseCallRecording') {
if (!this.srsClient || this.srsClient.paused) {
this.logger.info('discarding invalid pauseCallRecording request');
res.send(400);
return;
}
succeeded = await this.srsClient.pause();
}
else if (reason === 'resumeCallRecording') {
if (!this.srsClient || !this.srsClient.paused) {
res.send(400);
this.logger.info('discarding invalid resumeCallRecording request');
return;
}
succeeded = await this.srsClient.resume();
}
res.send(succeeded ? 200 : 503);
}
}
else if (dlg.type === 'uac' && ['application/dtmf-relay', 'application/dtmf'].includes(contentType)) {
const arr = /Signal=\s*([0-9#*])/.exec(req.body);
if (!arr) {
this.logger.info({body: req.body}, '_onInfo: invalid INFO dtmf request');
throw new Error(`_onInfo: no dtmf in body for ${contentType}`);
}
const code = arr[1];
const arr2 = /Duration=\s*(\d+)/.exec(req.body);
const duration = arr2 ? arr2[1] : 250;
if (this.isMediaReleased) {
/* just relay on to the feature server */
this.logger.info({code, duration}, 'got SIP INFO DTMF from caller, relaying to feature server');
this._onDTMF(dlg.other, {event: code, duration})
.catch((err) => this.logger.info({err}, 'Error relaying DTMF to feature server'));
res.send(200);
}
else {
/* else convert SIP INFO to RFC 2833 telephony events */
this.logger.info({code, duration}, 'got SIP INFO DTMF from caller, converting to RFC 2833');
const opts = {
...this.rtpEngineOpts.common,
'from-tag': this.rtpEngineOpts.uac.tag,
code,
duration
};
const response = await this.playDTMF(opts);
if ('ok' !== response.result) {
this.logger.info({response}, `rtpengine playDTMF failed with ${JSON.stringify(response)}`);
throw new Error('rtpengine failed: answer');
}
res.send(200);
}
}
else {
const response = await dlg.other.request({
@@ -610,10 +845,10 @@ Duration=${payload.duration} `
res.send(202);
// invite to new fs
const headers = {};
if (req.has('X-Retain-Call-Sid')) {
Object.assign(headers, {'X-Retain-Call-Sid': req.get('X-Retain-Call-Sid')});
}
const headers = {
...(req.has('X-Retain-Call-Sid') && {'X-Retain-Call-Sid': req.get('X-Retain-Call-Sid')}),
...(req.has('X-Account-Sid') && {'X-Account-Sid': req.get('X-Account-Sid')})
};
const dlg = await this.srf.createUAC(referTo.uri, {localSdp: dlg.local.sdp, headers});
this.uas = dlg;
this.uas.other = this.uac;
@@ -624,6 +859,7 @@ Duration=${payload.duration} `
this.logger.info('call ended with normal termination');
this.rtpEngineResource.destroy();
this.activeCallIds.delete(this.req.get('Call-ID'));
if (this.activeCallIds.size === 0) this.idleEmitter.emit('idle');
this.uas.other.destroy();
this.srf.endSession(this.req);
});
+84 -38
View File
@@ -1,7 +1,7 @@
const debug = require('debug')('jambonz:sbc-outbound');
const parseUri = require('drachtio-srf').parseUri;
const Registrar = require('@jambonz/mw-registrar');
const {selectHostPort, makeCallCountKey} = require('./utils');
const {selectHostPort, nudgeCallCounts} = require('./utils');
const FS_UUID_SET_NAME = 'fsUUIDs';
module.exports = (srf, logger, opts) => {
@@ -10,14 +10,21 @@ module.exports = (srf, logger, opts) => {
const registrar = new Registrar(opts);
const {
lookupAccountCapacitiesBySid,
lookupAccountBySid
lookupAccountBySid,
queryCallLimits
} = srf.locals.dbHelpers;
const initLocals = async(req, res, next) => {
req.locals = req.locals || {};
const callId = req.get('Call-ID');
req.locals.account_sid = req.get('X-Account-Sid');
req.locals.logger = logger.child({callId, account_sid: req.locals.account_sid});
req.locals.application_sid = req.get('X-Application-Sid');
const traceId = req.locals.trace_id = req.get('X-Trace-ID');
req.locals.logger = logger.child({
callId,
traceId,
account_sid:
req.locals.account_sid});
if (!req.locals.account_sid) {
logger.info('missing X-Account-Sid on outbound call');
@@ -58,7 +65,6 @@ module.exports = (srf, logger, opts) => {
}
}
stats.increment('sbc.invites', ['direction:outbound']);
req.on('cancel', () => {
@@ -71,6 +77,7 @@ module.exports = (srf, logger, opts) => {
try {
req.locals.account = await lookupAccountBySid(req.locals.account_sid);
req.locals.service_provider_sid = req.locals.account.service_provider_sid;
} catch (err) {
req.locals.logger.error({err}, `Error looking up account sid ${req.locals.account_sid}`);
res.send(500);
@@ -80,26 +87,27 @@ module.exports = (srf, logger, opts) => {
};
const checkLimits = async(req, res, next) => {
const {logger, account_sid} = req.locals;
const {writeAlerts, AlertType} = req.srf.locals;
const {logger, account_sid, service_provider_sid, application_sid} = req.locals;
const trackingOn = process.env.JAMBONES_TRACK_ACCOUNT_CALLS ||
process.env.JAMBONES_TRACK_SP_CALLS ||
process.env.JAMBONES_TRACK_APP_CALLS;
if (!process.env.JAMBONES_HOSTING && !trackingOn) {
logger.debug('tracking is off, skipping call limit checks');
return next(); // skip
}
const {writeCallCount, writeCallCountSP, writeCallCountApp, writeAlerts, AlertType} = req.srf.locals;
const key = makeCallCountKey(account_sid);
try {
/* increment the call count */
const calls = await incrKey(key);
debug(`checkLimits: call count is now ${calls}`);
/* decrement count if INVITE is later rejected */
res.once('end', ({status}) => {
res.once('end', async({status}) => {
if (status > 200) {
debug('checkLimits: decrementing call count due to rejection');
decrKey(key)
.then((count) => {
logger.debug({key}, `after rejection there are ${count} active calls for this account`);
debug({key}, `after rejection there are ${count} active calls for this account`);
return;
})
.catch((err) => logger.error({err}, 'checkLimits: decrKey err'));
nudgeCallCounts(logger, {
service_provider_sid,
account_sid,
application_sid
}, decrKey, {writeCallCountSP, writeCallCount, writeCallCountApp})
.catch((err) => logger.error(err, 'Error decrementing call counts'));
const tags = ['accepted:no', `sipStatus:${status}`];
stats.increment('sbc.originations', tags);
}
@@ -109,6 +117,13 @@ module.exports = (srf, logger, opts) => {
}
});
/* increment the call count */
const {callsSP, calls} = await nudgeCallCounts(logger, {
service_provider_sid,
account_sid,
application_sid
}, incrKey, {writeCallCountSP, writeCallCount, writeCallCountApp});
/* compare to account's limit, though avoid db hit when call count is low */
const minLimit = process.env.MIN_CALL_LIMIT ?
parseInt(process.env.MIN_CALL_LIMIT) :
@@ -117,23 +132,54 @@ module.exports = (srf, logger, opts) => {
const capacities = await lookupAccountCapacitiesBySid(account_sid);
const limit = capacities.find((c) => c.category == 'voice_call_session');
if (!limit) {
logger.debug('checkLimits: no call limits specified');
return next();
if (limit) {
const limit_sessions = limit.quantity;
if (calls > limit_sessions) {
logger.info({calls, limit_sessions}, 'checkLimits: limits exceeded');
writeAlerts({
alert_type: AlertType.ACCOUNT_CALL_LIMIT,
service_provider_sid,
account_sid,
count: limit_sessions
}).catch((err) => logger.info({err}, 'checkLimits: error writing alert'));
res.send(503, 'Maximum Calls In Progress');
return req.srf.endSession(req);
}
}
const limit_sessions = limit.quantity;
if (calls > limit_sessions) {
debug(`checkLimits: limits exceeded: call count ${calls}, limit ${limit_sessions}`);
logger.info({calls, limit_sessions}, 'checkLimits: limits exceeded');
writeAlerts({
alert_type: AlertType.CALL_LIMIT,
account_sid,
count: limit_sessions
}).catch((err) => logger.info({err}, 'checkLimits: error writing alert'));
res.send(503, 'Maximum Calls In Progress');
return req.srf.endSession(req);
else if (trackingOn) {
const {account_limit, sp_limit} = await queryCallLimits(service_provider_sid, account_sid);
if (process.env.JAMBONES_TRACK_ACCOUNT_CALLS && account_limit > 0 && calls > account_limit) {
logger.info({calls, account_limit}, 'checkLimits: account limits exceeded');
writeAlerts({
alert_type: AlertType.ACCOUNT_CALL_LIMIT,
service_provider_sid: service_provider_sid,
account_sid,
count: calls
}).catch((err) => logger.info({err}, 'checkLimits: error writing alert'));
res.send(503, 'Max Account Calls In Progress', {
headers: {
'X-Account-Sid': account_sid,
'X-Call-Limit': account_limit
}
});
return req.srf.endSession(req);
}
if (process.env.JAMBONES_TRACK_SP_CALLS && sp_limit > 0 && callsSP > sp_limit) {
logger.info({callsSP, sp_limit}, 'checkLimits: service provider limits exceeded');
writeAlerts({
alert_type: AlertType.SP_CALL_LIMIT,
service_provider_sid: service_provider_sid,
count: callsSP
}).catch((err) => logger.info({err}, 'checkLimits: error writing alert'));
res.send(503, 'Max Service Provider Calls In Progress', {
headers: {
'X-Service-Provider-Sid': service_provider_sid,
'X-Call-Limit': sp_limit
}
});
return req.srf.endSession(req);
}
}
next();
} catch (err) {
@@ -148,8 +194,8 @@ module.exports = (srf, logger, opts) => {
logger.info(`received outbound INVITE to ${req.uri} from server at ${req.server.hostport}`);
const uri = parseUri(req.uri);
const desiredRouting = req.get('X-Jambonz-Routing');
if (!uri || !uri.user || !uri.host) {
const validUri = uri && uri.user && uri.host;
if (['user', 'sip'].includes(desiredRouting) && !validUri) {
logger.info({uri: req.uri}, 'invalid request-uri on outbound call, rejecting');
res.send(400, {
headers: {
+102 -7
View File
@@ -4,12 +4,21 @@ const debug = require('debug')('jambonz:sbc-outbound');
function makeRtpEngineOpts(req, srcIsUsingSrtp, dstIsUsingSrtp, teams = false) {
const from = req.getParsedHeader('from');
const srtpOpts = teams ? srtpCharacteristics['teams'] : srtpCharacteristics['default'];
const dstOpts = dstIsUsingSrtp ? srtpOpts : rtpCharacteristics;
const srcOpts = srcIsUsingSrtp ? srtpOpts : rtpCharacteristics;
const rtpCopy = JSON.parse(JSON.stringify(rtpCharacteristics));
const srtpCopy = JSON.parse(JSON.stringify(srtpCharacteristics));
const srtpOpts = teams ? srtpCopy['teams'] : srtpCopy['default'];
const dstOpts = dstIsUsingSrtp ? srtpOpts : rtpCopy;
const srcOpts = srcIsUsingSrtp ? srtpOpts : rtpCopy;
/* webrtc clients (e.g. sipjs) send DMTF via SIP INFO */
if ((srcIsUsingSrtp || dstIsUsingSrtp) && !teams) {
dstOpts.flags.push('inject DTMF');
srcOpts.flags.push('inject DTMF');
}
const common = {
'call-id': req.get('Call-ID'),
'replace': ['origin', 'session-connection']
'replace': ['origin', 'session-connection'],
'record call': process.env.JAMBONES_RECORD_ALL_CALLS ? 'yes' : 'no'
};
return {
common,
@@ -71,7 +80,9 @@ const pingMsTeamsGateways = (logger, srf) => {
});
};
const makeCallCountKey = (sid) => `${sid}:outcalls`;
const makeAccountCallCountKey = (sid) => `outcalls:account:${sid}`;
const makeSPCallCountKey = (sid) => `outcalls:sp:${sid}`;
const makeAppCallCountKey = (sid) => `outcalls:app:${sid}`;
const equalsIgnoreOrder = (a, b) => {
if (a.length !== b.length) return false;
@@ -84,10 +95,94 @@ const equalsIgnoreOrder = (a, b) => {
return true;
};
const systemHealth = async(redisClient, ping, getCount) => {
await Promise.all([redisClient.ping(), ping()]);
return getCount();
};
const createHealthCheckApp = (port, logger) => {
const express = require('express');
const app = express();
app.use(express.urlencoded({ extended: true }));
app.use(express.json());
return new Promise((resolve) => {
app.listen(port, () => {
logger.info(`Health check server started at http://localhost:${port}`);
resolve(app);
});
});
};
const nudgeCallCounts = async(logger, sids, nudgeOperator, writers) => {
const {service_provider_sid, account_sid, application_sid} = sids;
const {writeCallCount, writeCallCountSP, writeCallCountApp} = writers;
const nudges = [];
const writes = [];
logger.debug(sids, 'nudgeCallCounts');
if (process.env.JAMBONES_TRACK_SP_CALLS) {
const key = makeSPCallCountKey(service_provider_sid);
nudges.push(nudgeOperator(key));
}
else {
nudges.push(() => Promise.resolve(null));
}
if (process.env.JAMBONES_TRACK_ACCOUNT_CALLS || process.env.JAMBONES_HOSTING) {
const key = makeAccountCallCountKey(account_sid);
nudges.push(nudgeOperator(key));
}
else {
nudges.push(() => Promise.resolve(null));
}
if (process.env.JAMBONES_TRACK_APP_CALLS && application_sid) {
const key = makeAppCallCountKey(application_sid);
nudges.push(nudgeOperator(key));
}
else {
nudges.push(() => Promise.resolve(null));
}
try {
const [callsSP, calls, callsApp] = await Promise.all(nudges);
logger.debug({
calls, callsSP, callsApp,
service_provider_sid, account_sid, application_sid}, 'call counts after adjustment');
if (process.env.JAMBONES_TRACK_SP_CALLS) {
writes.push(writeCallCountSP({service_provider_sid, calls_in_progress: callsSP}));
}
if (process.env.JAMBONES_TRACK_ACCOUNT_CALLS || process.env.JAMBONES_HOSTING) {
writes.push(writeCallCount({service_provider_sid, account_sid, calls_in_progress: calls}));
}
if (process.env.JAMBONES_TRACK_APP_CALLS && application_sid) {
writes.push(writeCallCountApp({service_provider_sid, account_sid, application_sid, calls_in_progress: callsApp}));
}
/* write the call counts to the database */
Promise.all(writes).catch((err) => logger.error({err}, 'Error writing call counts'));
return {callsSP, calls, callsApp};
} catch (err) {
logger.error(err, 'error incrementing call counts');
}
return {callsSP: null, calls: null, callsApp: null};
};
module.exports = {
makeRtpEngineOpts,
selectHostPort,
pingMsTeamsGateways,
makeCallCountKey,
equalsIgnoreOrder
makeAccountCallCountKey,
makeSPCallCountKey,
equalsIgnoreOrder,
systemHealth,
createHealthCheckApp,
nudgeCallCounts
};
+2670 -4261
View File
File diff suppressed because it is too large Load Diff
+17 -13
View File
@@ -1,6 +1,6 @@
{
"name": "sbc-outbound",
"version": "v0.7.3",
"version": "v0.8.2",
"main": "app.js",
"engines": {
"node": ">= 12.0.0"
@@ -22,28 +22,32 @@
"description": "jambonz session border controller application for outbound calls",
"scripts": {
"start": "node app",
"test": "NODE_ENV=test JAMBONZ_HOSTING=1 JAMBONES_NETWORK_CIDR=127.0.0.1/32 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_TIME_SERIES_HOST=127.0.0.1 JAMBONES_LOGLEVEL=error DRACHTIO_SECRET=cymru DRACHTIO_HOST=127.0.0.1 DRACHTIO_PORT=9060 JAMBONES_RTPENGINES=127.0.0.1:12222 node test/ ",
"test": "NODE_ENV=test HTTP_PORT=3050 JAMBONES_HOSTING=1 JAMBONES_NETWORK_CIDR=127.0.0.1/32 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_TIME_SERIES_HOST=127.0.0.1 JAMBONES_LOGLEVEL=error DRACHTIO_SECRET=cymru DRACHTIO_HOST=127.0.0.1 DRACHTIO_PORT=9060 JAMBONES_RTPENGINES=127.0.0.1:12222 node test/ ",
"coverage": "./node_modules/.bin/nyc --reporter html --report-dir ./coverage npm run test",
"jslint": "eslint app.js lib"
},
"dependencies": {
"@jambonz/db-helpers": "^0.6.16",
"@jambonz/db-helpers": "^0.7.4",
"@jambonz/realtimedb-helpers": "^0.7.0",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/mw-registrar": "0.2.1",
"@jambonz/realtimedb-helpers": "^0.4.24",
"@jambonz/rtpengine-utils": "^0.3.1",
"@jambonz/stats-collector": "^0.1.6",
"@jambonz/time-series": "^0.1.6",
"@jambonz/mw-registrar": "0.2.2",
"@jambonz/rtpengine-utils": "^0.4.3",
"@jambonz/siprec-client-utils": "^0.2.4",
"@jambonz/stats-collector": "^0.1.8",
"@jambonz/time-series": "^0.2.5",
"cidr-matcher": "^2.1.1",
"debug": "^4.3.3",
"debug": "^4.3.4",
"drachtio-fn-b2b-sugar": "^0.0.12",
"drachtio-srf": "^4.4.59",
"pino": "^7.4.1"
"drachtio-srf": "^4.5.21",
"express": "^4.18.1",
"pino": "^7.11.0",
"sdp-transform": "^2.14.1"
},
"devDependencies": {
"bent": "^7.3.12",
"eslint": "^7.32.0",
"eslint-plugin-promise": "^5.1.1",
"eslint-plugin-promise": "^5.2.0",
"nyc": "^15.1.0",
"tape": "^5.3.2"
"tape": "^5.5.3"
}
}
+4
View File
@@ -5,3 +5,7 @@ DRACHTIO_SECRET=cymru
JAMBONES_REDIS_HOST=172.39.0.11
JAMBONES_REDIS_PORT=6379
JAMBONES_LOGLEVEL=info
JAMBONES_MYSQL_HOST=172.39.0.2
JAMBONES_MYSQL_USER=jambones_test
JAMBONES_MYSQL_PASSWORD=jambones_test
JAMBONES_MYSQL_DATABASE=jambones_test
@@ -0,0 +1,91 @@
<?xml version="1.0" encoding="ISO-8859-1" ?>
<!DOCTYPE scenario SYSTEM "sipp.dtd">
<scenario name="UAC with media">
<send retrans="500">
<![CDATA[
INVITE sip:16173333456@sbc-sip:5060 SIP/2.0
Via: SIP/2.0/[transport] [local_ip]:[local_port];branch=[branch]
From: sipp <sip:sipp@[local_ip]:[local_port]>;tag=[pid]SIPpTag09[call_number]
To: <sip:16173333456@sbc-sip:5060>
Call-ID: [call_id]
CSeq: 1 INVITE
Contact: sip:sipp@[local_ip]:[local_port]
Max-Forwards: 70
X-Account-Sid: ed649e33-e771-403a-8c99-1780eabbc803
X-Call-Sid: fff49e33-e771-403a-8c99-1780eabbc803
X-Jambonz-Routing: phone
Subject: uac-pcap-carrier-success
Content-Type: application/sdp
Content-Length: [len]
v=0
o=user1 53655765 2353687637 IN IP[local_ip_type] [local_ip]
s=-
c=IN IP[local_ip_type] [local_ip]
t=0 0
m=audio [auto_media_port] RTP/AVP 8 101
a=rtpmap:8 PCMA/8000
a=rtpmap:101 telephone-event/8000
a=fmtp:101 0-11,16
]]>
</send>
<recv response="100" optional="true">
</recv>
<recv response="180" optional="true">
</recv>
<recv response="200" rtd="true" crlf="true">
</recv>
<send>
<![CDATA[
ACK sip:sip:+16173333456@sbc-sip:5060 SIP/2.0
Via: SIP/2.0/[transport] [local_ip]:[local_port];branch=[branch]
From: sipp <sip:sipp@[local_ip]:[local_port]>;tag=[pid]SIPpTag09[call_number]
To: <sip:sip:+16173333456@sbc-sip:5060>[peer_tag_param]
Call-ID: [call_id]
CSeq: 1 ACK
Subject: uac-pcap-carrier-success
Content-Length: 0
]]>
</send>
<!-- Play a pre-recorded PCAP file (RTP stream) -->
<nop>
<action>
<exec play_pcap_audio="pcap/g711a.pcap"/>
</action>
</nop>
<!-- Pause briefly -->
<pause milliseconds="2000"/>
<!-- The 'crlf' option inserts a blank line in the statistics report. -->
<send retrans="500">
<![CDATA[
BYE sip:sip:+16173333456@sbc-sip:5060 SIP/2.0
Via: SIP/2.0/[transport] [local_ip]:[local_port];branch=[branch]
From: sipp <sip:sipp@[local_ip]:[local_port]>;tag=[pid]SIPpTag09[call_number]
To: <sip:sip:+16173333456@sbc-sip:5060>[peer_tag_param]
Call-ID: [call_id]
CSeq: 2 BYE
Max-Forwards: 70
Subject: uac-pcap-carrier-success
Content-Length: 0
]]>
</send>
<recv response="200" crlf="true">
</recv>
</scenario>
+14 -4
View File
@@ -2,7 +2,8 @@ const test = require('tape');
const { output, sippUac } = require('./sipp')('test_sbc-outbound');
const {execSync} = require('child_process');
const debug = require('debug')('jambonz:sbc-outbound');
const consoleLogger = {error: console.error, info: console.log, debug: console.log};
const bent = require('bent');
const getJSON = bent('json');
process.on('unhandledRejection', (reason, p) => {
console.log('Unhandled Rejection at: Promise', p, 'reason:', reason);
@@ -29,6 +30,11 @@ test('sbc-outbound tests', async(t) => {
try {
await connect(srf);
let obj = await getJSON('http://127.0.0.1:3050/');
t.ok(obj.calls === 0, 'HTTP GET / works (current call count)')
obj = await getJSON('http://127.0.0.1:3050/system-health');
t.ok(obj.calls === 0, 'HTTP GET /system-health works (health check)')
/* call to unregistered user */
debug('successfully connected to drachtio server');
await sippUac('uac-pcap-device-404.xml');
@@ -36,7 +42,11 @@ test('sbc-outbound tests', async(t) => {
/* call to PSTN with no lcr configured */
await sippUac('uac-pcap-carrier-success.xml');
t.pass('successfully completed outbound call to configured sip trunk');
t.pass('successfully completed outbound call to sip trunk');
/* call to PSTN with request uri we see in kubernetes */
await sippUac('uac-pcap-carrier-success-k8s.xml');
t.pass('successfully completed outbound call to sip trunk (k8S req uri)');
// re-rack test data
execSync(`mysql -h 127.0.0.1 -u root --protocol=tcp -D jambones_test < ${__dirname}/db/jambones-sql.sql`);
@@ -74,10 +84,10 @@ test('sbc-outbound tests', async(t) => {
await sippUac('uac-pcap-carrier-fail-limits.xml');
t.pass('fails when max calls in progress');
await waitFor(10);
await waitFor(25);
const res = await queryCdrs({account_sid: 'ed649e33-e771-403a-8c99-1780eabbc803'});
//console.log(`cdrs: ${JSON.stringify(res)}`);
console.log(`cdrs: ${JSON.stringify(res)}`);
t.ok(res.total === 6, 'wrote 6 cdrs');
srf.disconnect();