Compare commits

...
75 Commits
Author SHA1 Message Date
Dave Horton acead419d5 bump version to 0.7.6 2022-08-26 20:09:15 +02:00
Dave Horton 911d208e0f when releasing media anchor from FS, if we are using SRTP on A leg we need to reinvite 2022-08-24 14:26:24 +02:00
xquanluu c824367086 feat: update time-series 0.11.12 (#45) 2022-08-19 16:17:55 +02:00
Dave Horton f39df615f2 update time-series 2022-08-19 09:59:08 +02:00
Dave Horton 58e80afc92 deps 2022-08-12 14:09:19 +02:00
Dave Horton 738f151066 Feature/sip info dtmf (#44)
* initial changes for handling SIP INFO from webrtc clients

* SIP INFO DTMF handling: if media has been released, just relay SIP INFO to FS otherwise transcode to RFC 2833
2022-08-11 14:33:01 +02:00
Dave Horton a94f25b0bd update to latest @jambonz/siprec-utils with fix for sdp version 2022-08-08 16:27:42 +02:00
Dave Horton a1542b161b update siprec-utils to increment sdp version on reinvite 2022-08-08 13:48:52 +02:00
Dave Horton 4e961491a6 Feature/siprec server (#42)
* changes to handle siprec invites

* retain xml for sending

* send multipart body for siprec

* fixup sdp

* possibly trim xml string

* fix prev commit

* fix prev commit

* still tweaking sdp

* tweaks

* tweaking

* bit of refactoring

* refactor siprec client stuff into a library

* add direction in siprec metadata

* fix
2022-08-05 10:27:50 +01:00
Dave Horton 4e5f7ae908 bugfix: MS Teams warm transfer (invite w/replaces) 2022-08-01 14:49:28 +01:00
Dave Horton 798e070127 Dockerfile: update base image 2022-07-28 12:53:31 +01:00
Dave Horton 2d350d4850 bugfix: regression - adding custom headers to refer caused preferred hostname on Refer-To to be lost 2022-07-27 11:31:29 +01:00
Dave Horton baad125924 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:26:17 +01:00
Dave Horton 73528f5ce2 bugfix #131: pass on custom headers in REFER 2022-07-18 14:58:12 +02:00
Snyk bot c8329a94f3 fix: upgrade drachtio-srf from 4.5.0 to 4.5.1 (#39)
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/be6a10dd-83a0-4fef-a6b9-7dea4a956a5a?utm_source=github&utm_medium=referral&page=upgrade-pr
2022-07-13 09:17:33 +02:00
Dave Horton dceeef5549 update to latest verify-aws-sns-signature 2022-07-06 18:27:24 +02:00
Paulo Tellesandp.souza b5551fffba improve dockerfile to fix snyk security issues (#38)
Co-authored-by: p.souza <p.souza@cognigy.com>
2022-07-06 18:17:20 +02:00
Dave Horton d2b5597571 initial support for siprec recording (#36)
* initial support for siprec recording

* handle pause/resume siprec recording
2022-06-23 16:23:09 -04:00
Dave Horton c9401ab3c8 update to azure 1.22.0 2022-06-11 16:22:26 -04:00
Dave Horton cc5c712a5b update deps 2022-06-11 11:40:44 -04:00
Dave Horton a923227e4a fix test scenario 2022-05-15 10:05:57 -04:00
Dave Horton 5f7df4d135 #32 - allow wildcard matches 2022-05-07 11:17:21 -04:00
Dave Horton 7dd4d4a045 Feature/healthcheck improvements (#30)
* health check tests mysql and redis connectivity

* health check tests mysql and redis connectivity

* minor
2022-04-12 15:46:20 -04:00
Dave Horton 7382bf6bd6 bump version 2022-04-06 08:18:08 -04:00
Dave Horton abfce38150 logging 2022-04-01 15:12:44 -04:00
Dave Horton d8035c978e bugfix: error writing call counts 2022-04-01 15:09:46 -04:00
Dave Horton 91cd677ad8 write call_counts time series data tracking inbound call counts by account 2022-04-01 13:08:25 -04:00
Dave Horton 17b53438e2 track account level calls if env JAMBONES_TRACK_ACCOUNT_CALLS is set 2022-04-01 06:51:44 -04:00
Dave Horton f1e2f192b7 Bugfix/mem leak (#29)
* wait for 200 OK to BYE before closing connection to drachtio

* minor logging
2022-03-27 21:50:09 -04:00
Dave Horton 0af8bcc348 remove dlg event handlers on destroy 2022-03-27 20:23:30 -04:00
Dave Horton e067125974 de-link Dialogs so they can be GC'ed at call end 2022-03-27 19:43:39 -04:00
Dave Horton 2dadde64f4 write otel trace_id to call history 2022-03-23 09:29:22 -04:00
Dave Horton d5a1337811 bugfix: connection to drachtio was not closed on non-success response to invite (#27) 2022-03-18 14:04:52 -04:00
Snyk bot a6571134cd fix: Dockerfile to reduce vulnerabilities (#26)
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:54:19 -04:00
Dave Horton e0e5c75496 bump version 2022-03-08 20:15:17 -05:00
Dave Horton c08e35c261 Feature/incoming refer (#25)
* handle incoming REFER and send on to the FS

* clarity
2022-03-05 15:22:11 -05:00
Dave Horton 883c63723c update to dbhelpers with support for searching root domains to identify account 2022-02-17 20:59:17 -05:00
Dave Horton 13f78bf8d9 added pre-commit hook for linting 2022-02-14 13:18:04 -05:00
Dave Horton 55056c1771 bump version 2022-02-09 15:43:40 -05:00
Dave Horton 87ec5f8e09 update to latest @jambonz/realtimedb-helpers with support for redis username / password auth 2022-02-09 15:13:54 -05:00
Dave Horton 317280befc update deps 2022-02-09 08:22:14 -05:00
Dave Horton 7b97a0a137 0.7.2 version 2022-01-28 09:12:35 -05:00
Dave Horton 75d8381ebb Feature/rtpengine locate by dns (#22)
* initial changes to use udp and dns in K8s for rtpengine ng

* github actions

* logging

* remove port from service name before dns lookup

* query rtpengine dns immediately on startup

* bugfix prev checkin

* use dns.lookup instead of resolve4 (k8s does not appear to consult search patterns in resolve.conf with the latter)

* return all rtpengine endpoints in dns lookup

* bugfix: polling rtpengine endpoints

* fix prev commit

* catch error on rtpengine failure

* update deps

* require node 14 in gh action build

* specify node version in gh actions
2022-01-23 18:03:11 -05:00
Dave Horton 7750fdc3f2 healthcheck only in k8s 2022-01-22 21:44:14 -05:00
Dave Horton fcf90916e1 update to rtpengine-utils that reconnects ws 2022-01-19 19:27:30 -05:00
Dave Horton 4a0aa79456 no need to reset rtpengine session when releasing media, and it was causing ICE to be reset which caused Teams to drop 2022-01-17 18:39:09 -05:00
Dave Horton 1cdaae7773 use ws rather than tcp for rtpengine connection in K8S 2022-01-12 08:28:41 -05:00
Dave Horton 82ac86389d add support for tcp connections to rtpengine, needed for K8S (#21)
* add support for tcp connections to rtpengine, needed for K8S

* update to correct version of @jambonz/rtpengine-utils

* update deps
2022-01-11 15:42:49 -05:00
Dave Horton e9184ad208 add support for using ws to connect to rtpengine 2022-01-10 22:00:04 -05:00
Dave Horton 804bb890b5 JAMBONES_NETWORK_CIDR not needed for K8S (#20) 2022-01-09 14:59:00 -05:00
Dave Horton a892a87eb5 K8s changes (#19)
* K8S changes

* k8s: test explicit dns lookup of service

* bugfix prev commit

* typo

* k8s: more dns

* k8s: more dns

* k8s: more dns fun

* k8s cleanup

* k8s: user service for rtpengine location

* k8s cleanup

* typo

* change env name for fs in k8s

* change k8s service name for feature server

* add support for outbound connection mode

* k8s change for outbound

* minor

* bugfix: drachtio connection was dropped after successful connect

* drop drachtio connection on call end

* Dockerfile

* k8s pre-stop hook

* actual hook committed

* make hjook executable

* dockerfile change

* time series fix

* bugfix: teams transfer using replaces
2022-01-06 12:37:49 -05:00
Dave Horton 56205dc852 bump version 2021-12-21 09:41:01 -05:00
Dave Horton 927b7be637 SIGTERM handler to remove entry from active-sip 2021-12-20 16:09:34 -05:00
Dave Horton 63b482e562 bugfix: special case of single-tenant system 2021-12-20 12:23:04 -05:00
Dave Horton a7d047a7e8 remove unneeded refs 2021-12-20 10:09:51 -05:00
Dave Horton ea4f5ea0a8 bugfix: incoming call to accounts.sip_realm did not look for SP-level carriers 2021-12-20 09:54:21 -05:00
Dave Horton 13d8ee8c3c version bump, add docker publish 2021-12-13 09:52:41 -05:00
Dave Horton c717470c5e version bump 2021-12-02 19:20:57 -05:00
Dave Horton f16598e144 Feature/sip refer (#18)
* support for sip refer to transfer an incoming call

* when sending REFER for call transfer, format Refer-To with carrier trunk if applicable

* better handling of e164 on Refer-To header
2021-11-20 11:41:12 -05:00
Dave Horton 8584020d4c add support for proxies that add X-Forwarded-For 2021-11-05 09:29:12 -04:00
Dave Horton 0b1b49bd59 Dockerfile 2021-11-04 12:57:15 -04:00
Dave Horton e22b2ae63f version bump 2021-11-03 13:51:50 -04:00
Dave Horton 2a5200673d bump version 2021-10-21 13:08:53 -04:00
Dave Horton 13ca308df7 bump version 2021-10-21 13:01:04 -04:00
Dave Horton 981788dc58 Feature/minimal media anchoring (#11)
* add support for sitting behind a sip proxy that adds X-Forwarded-For header

* feature: release media from freeswitch

* handle mute/unmute

* fix for relaying INFO

* support for relaying dtmf via SIP INFO to FS

* deps
2021-10-21 12:00:02 -04:00
Dave Horton ca4599da22 bugfix: autoscaling 2021-10-02 17:51:04 -04:00
Dave Horton 04c4ad5975 add support for AWS autoscaling (#10) 2021-10-02 12:41:52 -04:00
Dave Horton 5267b6f9a4 deps 2021-10-01 15:37:41 -04:00
Dave Horton 2fa1f08244 slight change to data format written to redis for set of active sbcs 2021-10-01 15:37:07 -04:00
Dave Horton 03cb1cceb5 maintain a redis set of active SBC SIP servers, under {prefix}:active-sip 2021-09-29 18:28:37 -04:00
Dave Horton 9b78a818b5 bugfix: reinvite handling mixed up the public-private direction in rtpengine offer 2021-08-09 15:44:23 -04:00
Dave Horton e6a90a32ad if we get re-invite with no SDP (looking at you, BT) just respond with current offer 2021-08-03 10:41:41 -04:00
Dave Horton fb21fd2333 brackets around From header 2021-07-28 16:16:47 -04:00
Dave Horton 819fecf455 LICENSE 2021-07-21 12:40:25 -04:00
Dave Horton b0406287f5 fix tests 2021-06-29 12:47:30 -04:00
27 changed files with 4535 additions and 1267 deletions
+1 -1
View File
@@ -8,7 +8,7 @@
"jsx": false,
"modules": false
},
"ecmaVersion": 2018
"ecmaVersion": 2020
},
"plugins": ["promise"],
"rules": {
+51
View File
@@ -0,0 +1,51 @@
name: Docker
on:
push:
# Publish `master` as Docker `latest` image.
branches:
- main
# Publish `v1.2.3` tags as releases.
tags:
- v*
env:
IMAGE_NAME: sbc-inbound
jobs:
push:
runs-on: ubuntu-latest
if: github.event_name == 'push'
steps:
- uses: actions/checkout@v2
- name: Build image
run: docker build . --file Dockerfile --tag $IMAGE_NAME
- name: Log into registry
run: echo "${{ secrets.GITHUB_TOKEN }}" | docker login ghcr.io -u ${{ github.actor }} --password-stdin
- name: Push image
run: |
IMAGE_ID=ghcr.io/${{ github.repository_owner }}/$IMAGE_NAME
# Change all uppercase to lowercase
IMAGE_ID=$(echo $IMAGE_ID | tr '[A-Z]' '[a-z]')
# Strip git ref prefix from version
VERSION=$(echo "${{ github.ref }}" | sed -e 's,.*/\(.*\),\1,')
# Strip "v" prefix from tag name
[[ "${{ github.ref }}" == "refs/tags/"* ]] && VERSION=$(echo $VERSION | sed -e 's/^v//')
# Use Docker `latest` tag convention
[ "$VERSION" == "main" ] && VERSION=latest
echo IMAGE_ID=$IMAGE_ID
echo VERSION=$VERSION
docker tag $IMAGE_NAME $IMAGE_ID:$VERSION
docker push $IMAGE_ID:$VERSION
@@ -11,8 +11,8 @@ jobs:
- uses: actions/checkout@v2
- uses: actions/setup-node@v1
with:
node-version: 12
- run: npm install
node-version: 14.x
- run: npm ci
- run: npm run jslint
- run: npm test
+1 -1
View File
@@ -1,7 +1,7 @@
# Logs
logs
*.log
.vscode/
# Runtime data
pids
*.pid
+4
View File
@@ -0,0 +1,4 @@
#!/bin/sh
. "$(dirname "$0")/_/husky.sh"
npm run jslint
+18 -11
View File
@@ -1,16 +1,23 @@
FROM node:alpine as builder
RUN apk update && apk add --no-cache python make g++
WORKDIR /opt/app/
COPY package.json ./
RUN npm install
RUN npm prune
FROM --platform=linux/amd64 node:18.6.0-alpine as base
FROM node:alpine as app
WORKDIR /opt/app
COPY . /opt/app
COPY --from=builder /opt/app/node_modules ./node_modules
RUN apk --update --no-cache add --virtual .builds-deps build-base python3
WORKDIR /opt/app/
FROM base as build
COPY package.json package-lock.json ./
RUN npm ci
COPY . .
FROM base
COPY --from=build /opt/app /opt/app/
ARG NODE_ENV
ENV NODE_ENV $NODE_ENV
CMD [ "npm", "start" ]
CMD [ "node", "app.js" ]
+1 -1
View File
@@ -1,6 +1,6 @@
MIT License
Copyright (c) 2019 jambonz
Copyright (c) 2021 Drachtio Communications Services, LLC
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
+144 -28
View File
@@ -6,32 +6,36 @@ assert.ok(process.env.JAMBONES_MYSQL_HOST &&
assert.ok(process.env.DRACHTIO_PORT || process.env.DRACHTIO_HOST, 'missing DRACHTIO_PORT env var');
assert.ok(process.env.DRACHTIO_SECRET, 'missing DRACHTIO_SECRET env var');
assert.ok(process.env.JAMBONES_TIME_SERIES_HOST, 'missing JAMBONES_TIME_SERIES_HOST env var');
assert.ok(process.env.JAMBONES_NETWORK_CIDR, 'missing JAMBONES_NETWORK_CIDR env var');
assert.ok(process.env.JAMBONES_NETWORK_CIDR || process.env.K8S, 'missing JAMBONES_NETWORK_CIDR env var');
const Srf = require('drachtio-srf');
const srf = new Srf('sbc-inbound');
const CIDRMatcher = require('cidr-matcher');
const matcher = new CIDRMatcher([process.env.JAMBONES_NETWORK_CIDR]);
const opts = Object.assign({
timestamp: () => {return `, "time": "${new Date().toISOString()}"`;}
}, {level: process.env.JAMBONES_LOGLEVEL || 'info'});
const logger = require('pino')(opts);
const {
writeCallCount,
queryCdrs,
writeCdrs,
writeAlerts,
AlertType
} = require('@jambonz/time-series')(logger, {
host: process.env.JAMBONES_TIME_SERIES_HOST,
port: process.env.JAMBONES_TIME_SERIES_PORT || 8086,
commitSize: 50,
commitInterval: 'test' === process.env.NODE_ENV ? 7 : 20
});
const StatsCollector = require('@jambonz/stats-collector');
const stats = new StatsCollector(logger);
const {equalsIgnoreOrder, createHealthCheckApp, systemHealth} = require('./lib/utils');
const {LifeCycleEvents} = require('./lib/constants');
const setNameRtp = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-rtp`;
const rtpServers = [];
const setName = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-sip`;
const {
pool,
ping,
lookupAuthHook,
lookupSipGatewayBySignalingAddress,
addSbcAddress,
@@ -47,14 +51,26 @@ const {
database: process.env.JAMBONES_MYSQL_DATABASE,
connectionLimit: process.env.JAMBONES_MYSQL_CONNECTION_LIMIT || 10
}, logger);
const {createSet, retrieveSet, incrKey, decrKey} = require('@jambonz/realtimedb-helpers')({
const {
client: redisClient,
createSet,
retrieveSet,
addToSet,
removeFromSet,
incrKey,
decrKey} = require('@jambonz/realtimedb-helpers')({
host: process.env.JAMBONES_REDIS_HOST || 'localhost',
port: process.env.JAMBONES_REDIS_PORT || 6379
}, logger);
const {getRtpEngine, setRtpEngines} = require('@jambonz/rtpengine-utils')([], logger, {emitter: stats});
const {getRtpEngine, setRtpEngines} = require('@jambonz/rtpengine-utils')([], logger, {
emitter: stats,
dtmfListenPort: process.env.DTMF_LISTEN_PORT || 22224,
protocol: 'udp'
});
srf.locals = {...srf.locals,
stats,
writeCallCount,
queryCdrs,
writeCdrs,
writeAlerts,
@@ -63,6 +79,7 @@ srf.locals = {...srf.locals,
getRtpEngine,
dbHelpers: {
pool,
ping,
lookupAuthHook,
lookupSipGatewayBySignalingAddress,
lookupAccountByPhoneNumber,
@@ -78,36 +95,64 @@ srf.locals = {...srf.locals,
retrieveSet
}
};
srf.locals.getFeatureServer = require('./lib/fs-tracking')(srf, logger);
const {
wasOriginatedFromCarrier,
getApplicationForDidAndCarrier,
getOutboundGatewayForRefer
} = require('./lib/db-utils')(srf, logger);
srf.locals = {
...srf.locals,
wasOriginatedFromCarrier,
getApplicationForDidAndCarrier,
getOutboundGatewayForRefer,
getFeatureServer: require('./lib/fs-tracking')(srf, logger)
};
const activeCallIds = srf.locals.activeCallIds;
const {
initLocals,
handleSipRec,
identifyAccount,
checkLimits,
challengeDeviceCalls
} = require('./lib/middleware')(srf, logger);
const CallSession = require('./lib/call-session');
if (process.env.DRACHTIO_HOST) {
if (process.env.DRACHTIO_HOST && !process.env.K8S) {
const CIDRMatcher = require('cidr-matcher');
const cidrs = process.env.JAMBONES_NETWORK_CIDR
.split(',')
.map((s) => s.trim());
const matcher = new CIDRMatcher(cidrs);
srf.connect({host: process.env.DRACHTIO_HOST, port: process.env.DRACHTIO_PORT, secret: process.env.DRACHTIO_SECRET });
srf.on('connect', (err, hp) => {
if (err) return this.logger.error({err}, 'Error connecting to drachtio server');
logger.info(`connected to drachtio listening on ${hp}`);
if (process.env.SBC_ACCOUNT_SID) return;
const hostports = hp.split(',');
for (const hp of hostports) {
const arr = /^(.*)\/(.*):(\d+)$/.exec(hp);
if (arr && 'udp' === arr[1] && !matcher.contains(arr[2])) {
logger.info(`adding sbc address ${arr[2]}`);
logger.info(`adding sbc public address to database: ${arr[2]}`);
srf.locals.sipAddress = arr[2];
addSbcAddress(arr[2]);
if (!process.env.SBC_ACCOUNT_SID) addSbcAddress(arr[2]);
}
else if (arr && 'tcp' === arr[1] && matcher.contains(arr[2])) {
const hostport = `${arr[2]}:${arr[3]}`;
logger.info(`adding sbc private address to redis: ${hostport}`);
srf.locals.privateSipAddress = hostport;
srf.locals.addToRedis = () => addToSet(setName, hostport);
srf.locals.removeFromRedis = () => removeFromSet(setName, hostport);
srf.locals.addToRedis();
}
}
});
}
else {
srf.on('listening', () => {
logger.info(`listening in outbound mode on port ${process.env.DRACHTIO_PORT}`);
});
srf.listen({port: process.env.DRACHTIO_PORT, secret: process.env.DRACHTIO_SECRET});
}
if (process.env.NODE_ENV === 'test') {
@@ -117,7 +162,13 @@ if (process.env.NODE_ENV === 'test') {
}
/* install middleware */
srf.use('invite', [initLocals, identifyAccount, checkLimits, challengeDeviceCalls]);
srf.use('invite', [
initLocals,
handleSipRec,
identifyAccount,
checkLimits,
challengeDeviceCalls
]);
srf.invite((req, res) => {
if (req.has('Replaces')) {
@@ -140,33 +191,78 @@ srf.use((req, res, next, err) => {
res.send(500);
});
/* update call stats periodically */
setInterval(() => {
stats.gauge('sbc.sip.calls.count', activeCallIds.size, ['direction:inbound']);
}, 20000);
if (process.env.K8S || process.env.HTTP_PORT) {
const PORT = process.env.HTTP_PORT || 3000;
const healthCheck = require('@jambonz/http-health-check');
const arrayCompare = (a, b) => {
if (a.length !== b.length) return false;
const uniqueValues = new Set([...a, ...b]);
for (const v of uniqueValues) {
const aCount = a.filter((e) => e === v).length;
const bCount = b.filter((e) => e === v).length;
if (aCount !== bCount) return false;
}
return true;
const getCount = () => srf.locals.activeCallIds.size;
createHealthCheckApp(PORT, logger)
.then((app) => {
healthCheck({
app,
logger,
path: '/',
fn: getCount
});
healthCheck({
app,
logger,
path: '/system-health',
fn: systemHealth.bind(null, redisClient, ping, getCount)
});
return;
})
.catch((err) => {
logger.error({err}, 'Error creating health check server');
});
}
if ('test' !== process.env.NODE_ENV) {
/* update call stats periodically */
setInterval(() => {
stats.gauge('sbc.sip.calls.count', activeCallIds.size, ['direction:inbound']);
}, 20000);
}
const lookupRtpServiceEndpoints = (lookup, serviceName) => {
logger.debug(`dns lookup for ${serviceName}..`);
lookup(serviceName, {family: 4, all: true}, (err, addresses) => {
if (err) {
logger.error({err}, `Error looking up ${serviceName}`);
return;
}
logger.debug({addresses, rtpServers}, `dns lookup for ${serviceName} returned`);
const addrs = addresses.map((a) => a.address);
if (!equalsIgnoreOrder(addrs, rtpServers)) {
rtpServers.length = 0;
Array.prototype.push.apply(rtpServers, addrs);
logger.info({rtpServers}, 'rtpserver endpoints have been updated');
setRtpEngines(rtpServers.map((a) => `${a}:${process.env.RTPENGINE_PORT || 22222}`));
}
});
};
/* update rtpengines periodically */
if (process.env.JAMBONES_RTPENGINES) {
if (process.env.K8S_RTPENGINE_SERVICE_NAME) {
/* poll dns for endpoints every so often */
const arr = /^(.*):(\d+)$/.exec(process.env.K8S_RTPENGINE_SERVICE_NAME);
const svc = arr[1];
logger.info(`rtpengine(s) will be found at dns name: ${svc}`);
const {lookup} = require('dns');
lookupRtpServiceEndpoints(lookup, svc);
setInterval(lookupRtpServiceEndpoints.bind(null, lookup, svc), process.env.RTPENGINE_DNS_POLL_INTERVAL || 10000);
}
else if (process.env.JAMBONES_RTPENGINES) {
/* static list of rtpengines */
setRtpEngines([process.env.JAMBONES_RTPENGINES]);
}
else {
/* poll redis periodically for rtpengines that have registered via OPTIONS ping */
const getActiveRtpServers = async() => {
try {
const set = await retrieveSet(setNameRtp);
const newArray = Array.from(set);
logger.debug({newArray, rtpServers}, 'getActiveRtpServers');
if (!arrayCompare(newArray, rtpServers)) {
if (!equalsIgnoreOrder(newArray, rtpServers)) {
logger.info({newArray}, 'resetting active rtpengines');
setRtpEngines(newArray.map((a) => `${a}:${process.env.RTPENGINE_PORT || 22222}`));
rtpServers.length = 0;
@@ -176,11 +272,31 @@ else {
logger.error({err}, 'Error setting new rtpengines');
}
};
setInterval(() => {
getActiveRtpServers();
}, 30000);
getActiveRtpServers();
}
const {lifecycleEmitter} = require('./lib/autoscale-manager')(logger);
/* if we are scaling in, check every so often if call count has gone to zero */
setInterval(async() => {
if (lifecycleEmitter.operationalState === LifeCycleEvents.ScaleIn) {
if (0 === activeCallIds.size) {
logger.info('scale-in complete now that calls have dried up');
lifecycleEmitter.scaleIn();
}
}
}, 20000);
process.on('SIGUSR2', handle.bind(null, removeFromSet, setName));
process.on('SIGTERM', handle.bind(null, removeFromSet, setName));
function handle(removeFromSet, setName, signal) {
logger.info(`got signal ${signal}, removing ${srf.locals.privateSipAddress} from set ${setName}`);
removeFromSet(setName, srf.locals.privateSipAddress);
}
module.exports = {srf, logger};
+29
View File
@@ -0,0 +1,29 @@
#!/usr/bin/env node
const bent = require('bent');
const getJSON = bent('json');
const PORT = process.env.HTTP_PORT || 3000;
const sleep = (ms) => {
return new Promise((resolve) => setTimeout(resolve, ms));
};
(async function() {
try {
do {
const obj = await getJSON(`http://127.0.0.1:${PORT}/`);
const {calls} = obj;
if (calls === 0) {
console.log('no calls on the system, we can exit');
process.exit(0);
}
else {
console.log(`waiting for ${calls} to exit..`);
}
await sleep(10000);
} while (1);
} catch (err) {
console.error(err, 'Error querying health endpoint');
process.exit(-1);
}
})();
+1 -1
View File
@@ -3,6 +3,6 @@
"DTLS": "off",
"SDES": "off",
"ICE": "remove",
"flags": ["media handover"],
"flags": ["media handover", "port latching"],
"rtcp-mux": ["accept"]
}
+2 -2
View File
@@ -3,14 +3,14 @@
"transport-protocol": "UDP/TLS/RTP/SAVPF",
"ICE": "force",
"SDES": "off",
"flags": ["generate mid", "SDES-no", "media handover"],
"flags": ["generate mid", "SDES-no", "media handover", "port latching"],
"rtcp-mux": ["require"]
},
"teams": {
"transport-protocol": "RTP/SAVP",
"ICE": "force",
"SDES": "off",
"flags": ["generate mid", "SDES-no", "media handover"],
"flags": ["generate mid", "SDES-no", "media handover", "port latching"],
"rtcp-mux": ["accept"]
}
}
+63
View File
@@ -0,0 +1,63 @@
const noopLogger = {info: () => {}, error: () => {}};
const {LifeCycleEvents} = require('./constants');
const Emitter = require('events');
module.exports = (logger) => {
logger = logger || noopLogger;
// listen for SNS lifecycle changes
let lifecycleEmitter = new Emitter();
lifecycleEmitter.dryUpCalls = false;
if (process.env.AWS_SNS_TOPIC_ARM) {
(async function() {
try {
lifecycleEmitter = await require('./aws-sns-lifecycle')(logger);
lifecycleEmitter
.on(LifeCycleEvents.ScaleIn, async() => {
logger.info('AWS scale-in notification: begin drying up calls');
lifecycleEmitter.dryUpCalls = true;
lifecycleEmitter.operationalState = LifeCycleEvents.ScaleIn;
const {srf} = require('..');
const {activeCallIds, removeFromRedis} = srf.locals;
/* remove our private IP from the set of active SBCs so rtp and fs know we are gone */
removeFromRedis();
/* if we have zero calls, we can complete the scale-in right now */
const calls = activeCallIds.size;
if (0 === calls) {
logger.info('scale-in can complete immediately as we have no calls in progress');
lifecycleEmitter.completeScaleIn();
}
else {
logger.info(`${calls} calls in progress; scale-in will complete when they are done`);
}
})
.on(LifeCycleEvents.StandbyEnter, () => {
lifecycleEmitter.dryUpCalls = true;
const {srf} = require('..');
const {removeFromRedis} = srf.locals;
removeFromRedis();
logger.info('AWS enter pending state notification: begin drying up calls');
})
.on(LifeCycleEvents.StandbyExit, () => {
lifecycleEmitter.dryUpCalls = false;
const {srf} = require('..');
const {addToRedis} = srf.locals;
addToRedis();
logger.info('AWS exit pending state notification: re-enable calls');
});
} catch (err) {
logger.error({err}, 'Failure creating SNS notifier, lifecycle events will be disabled');
}
})();
}
return {lifecycleEmitter};
};
+186
View File
@@ -0,0 +1,186 @@
const Emitter = require('events');
const bent = require('bent');
const assert = require('assert');
const PORT = process.env.AWS_SNS_PORT || 3001;
const {LifeCycleEvents} = require('./constants');
const express = require('express');
const app = express();
const getString = bent('string');
const AWS = require('aws-sdk');
const sns = new AWS.SNS({apiVersion: '2010-03-31'});
const autoscaling = new AWS.AutoScaling({apiVersion: '2011-01-01'});
const {Parser} = require('xml2js');
const parser = new Parser();
const {validatePayload} = require('verify-aws-sns-signature');
AWS.config.update({region: process.env.AWS_REGION});
class SnsNotifier extends Emitter {
constructor(logger) {
super();
this.logger = logger;
}
async _handlePost(req, res) {
try {
const parsedBody = JSON.parse(req.body);
this.logger.debug({headers: req.headers, body: parsedBody}, 'Received HTTP POST from AWS');
if (!validatePayload(parsedBody)) {
this.logger.info('incoming AWS SNS HTTP POST failed signature validation');
return res.sendStatus(403);
}
this.logger.debug('incoming HTTP POST passed validation');
res.sendStatus(200);
switch (parsedBody.Type) {
case 'SubscriptionConfirmation':
const response = await getString(parsedBody.SubscribeURL);
const result = await parser.parseStringPromise(response);
this.subscriptionArn = result.ConfirmSubscriptionResponse.ConfirmSubscriptionResult[0].SubscriptionArn[0];
this.subscriptionRequestId = result.ConfirmSubscriptionResponse.ResponseMetadata[0].RequestId[0];
this.logger.info({
subscriptionArn: this.subscriptionArn,
subscriptionRequestId: this.subscriptionRequestId
}, 'response from SNS SubscribeURL');
const data = await this.describeInstance();
this.lifecycleState = data.AutoScalingInstances[0].LifecycleState;
break;
case 'Notification':
if (parsedBody.Subject.startsWith('Auto Scaling: Lifecycle action \'TERMINATING\'')) {
const msg = JSON.parse(parsedBody.Message);
if (msg.EC2InstanceId === this.instanceId) {
this.logger.info('SnsNotifier - begin scale-in operation');
this.scaleInParams = {
AutoScalingGroupName: msg.AutoScalingGroupName,
LifecycleActionResult: 'CONTINUE',
LifecycleActionToken: msg.LifecycleActionToken,
LifecycleHookName: msg.LifecycleHookName
};
this.operationalState = LifeCycleEvents.ScaleIn;
this.emit(LifeCycleEvents.ScaleIn);
this.unsubscribe();
}
else {
this.logger.debug(`SnsNotifier - instance ${msg.EC2InstanceId} is scaling in (not us)`);
}
}
break;
default:
this.logger.info(`unhandled SNS Post Type: ${parsedBody.Type}`);
}
} catch (err) {
this.logger.error({err}, 'Error processing SNS POST request');
if (!res.headersSent) res.sendStatus(500);
}
}
async init() {
try {
this.logger.info('SnsNotifier: retrieving instance data');
this.instanceId = await getString('http://169.254.169.254/latest/meta-data/instance-id');
this.publicIp = await getString('http://169.254.169.254/latest/meta-data/public-ipv4');
this.snsEndpoint = `http://${this.publicIp}:${PORT}`;
this.logger.info({
instanceId: this.instanceId,
publicIp: this.publicIp,
snsEndpoint: this.snsEndpoint
}, 'retrieved AWS instance data');
// start listening
app.use(express.urlencoded({ extended: true }));
app.use(express.json());
app.use(express.text());
app.post('/', this._handlePost.bind(this));
app.use((err, req, res, next) => {
this.logger.error(err, 'burped error');
res.status(err.status || 500).json({msg: err.message});
});
app.listen(PORT);
} catch (err) {
this.logger.error({err}, 'Error retrieving AWS instance metadata');
}
}
async subscribe() {
try {
const response = await sns.subscribe({
Protocol: 'http',
TopicArn: process.env.AWS_SNS_TOPIC_ARM,
Endpoint: this.snsEndpoint
}).promise();
this.logger.info({response}, `response to SNS subscribe to ${process.env.AWS_SNS_TOPIC_ARM}`);
} catch (err) {
this.logger.error({err}, `Error subscribing to SNS topic arn ${process.env.AWS_SNS_TOPIC_ARM}`);
}
}
async unsubscribe() {
if (!this.subscriptionArn) throw new Error('SnsNotifier#unsubscribe called without an active subscription');
try {
const response = await sns.unsubscribe({
SubscriptionArn: this.subscriptionArn
}).promise();
this.logger.info({response}, `response to SNS unsubscribe to ${process.env.AWS_SNS_TOPIC_ARM}`);
} catch (err) {
this.logger.error({err}, `Error unsubscribing to SNS topic arn ${process.env.AWS_SNS_TOPIC_ARM}`);
}
}
completeScaleIn() {
assert(this.scaleInParams);
autoscaling.completeLifecycleAction(this.scaleInParams, (err, response) => {
if (err) return this.logger.error({err}, 'Error completing scale-in');
this.logger.info({response}, 'Successfully completed scale-in action');
});
}
describeInstance() {
return new Promise((resolve, reject) => {
if (!this.instanceId) return reject('instance-id unknown');
autoscaling.describeAutoScalingInstances({
InstanceIds: [this.instanceId]
}, (err, data) => {
if (err) {
this.logger.error({err}, 'Error describing instances');
reject(err);
} else {
this.logger.info({data}, 'SnsNotifier: describeInstance');
resolve(data);
}
});
});
}
}
module.exports = async function(logger) {
const notifier = new SnsNotifier(logger);
await notifier.init();
await notifier.subscribe();
process.on('SIGHUP', async() => {
try {
const data = await notifier.describeInstance();
const state = data.AutoScalingInstances[0].LifecycleState;
if (state !== notifier.lifecycleState) {
notifier.lifecycleState = state;
switch (state) {
case 'Standby':
notifier.emit(LifeCycleEvents.StandbyEnter);
break;
case 'InService':
notifier.emit(LifeCycleEvents.StandbyExit);
break;
}
}
} catch (err) {
console.error(err);
}
});
return notifier;
};
+439 -37
View File
@@ -1,7 +1,8 @@
const Emitter = require('events');
const SrsClient = require('@jambonz/siprec-client-utils');
const {makeRtpEngineOpts, SdpWantsSrtp, makeCallCountKey} = require('./utils');
const {forwardInDialogRequests} = require('drachtio-fn-b2b-sugar');
const {parseUri, SipError} = require('drachtio-srf');
const {parseUri, stringifyUri, SipError} = require('drachtio-srf');
const debug = require('debug')('jambonz:sbc-inbound');
const MS_TEAMS_USER_AGENT = 'Microsoft.PSTNHub.SIPProxy';
const MS_TEAMS_SIP_ENDPOINT = 'sip.pstnhub.microsoft.com';
@@ -13,8 +14,22 @@ const MS_TEAMS_SIP_ENDPOINT = 'sip.pstnhub.microsoft.com';
const createBLegFromHeader = (req) => {
const from = req.getParsedHeader('From');
const uri = parseUri(from.uri);
if (uri && uri.user) return `sip:${uri.user}@localhost`;
return 'sip:anonymous@localhost';
if (uri && uri.user) return `<sip:${uri.user}@localhost>`;
return '<sip:anonymous@localhost>';
};
const createSiprecBody = (headers, sdp, type, content) => {
const sep = 'uniqueBoundary';
headers['Content-Type'] = `multipart/mixed;boundary="${sep}"`;
return `--${sep}\r
Content-Type: application/sdp\r
\r
${sdp}\r
--${sep}\r
Content-Type: ${type}\r
Content-Disposition: recording-session\r
\r
${content}`;
};
class CallSession extends Emitter {
@@ -24,6 +39,8 @@ class CallSession extends Emitter {
this.res = res;
this.srf = req.srf;
this.logger = logger.child({callId: req.get('Call-ID')});
this.siprec = req.locals.siprec;
this.xml = req.locals.xml;
this.getRtpEngine = req.srf.locals.getRtpEngine;
this.getFeatureServer = req.srf.locals.getFeatureServer;
@@ -33,13 +50,28 @@ class CallSession extends Emitter {
this.decrKey = req.srf.locals.realtimeDbHelpers.decrKey;
this.callCountKey = makeCallCountKey(req.locals.account_sid);
this._mediaReleased = false;
}
get isFromMSTeams() {
return !!this.req.locals.msTeamsTenantFqdn;
}
get privateSipAddress() {
return this.srf.locals.privateSipAddress;
}
get isMediaReleased() {
return this._mediaReleased;
}
get callerIsUsingSrtp() {
const tp = this.rtpEngineOpts?.uas?.mediaOpts['transport-protocol'];
return tp && -1 !== tp.indexOf('SAVP');
}
async connect() {
const {sdp} = this.req.locals;
this.logger.info('inbound call accepted for routing');
const engine = this.getRtpEngine();
if (!engine) {
@@ -49,10 +81,34 @@ class CallSession extends Emitter {
return this.res.send(480);
}
debug(`got engine: ${JSON.stringify(engine)}`);
const {offer, answer, del} = engine;
const {
offer,
answer,
del,
blockMedia,
unblockMedia,
blockDTMF,
unblockDTMF,
playDTMF,
subscribeDTMF,
unsubscribeDTMF,
subscribeRequest,
subscribeAnswer,
unsubscribe
} = engine;
this.offer = offer;
this.answer = answer;
this.del = del;
this.blockMedia = blockMedia;
this.unblockMedia = unblockMedia;
this.blockDTMF = blockDTMF;
this.unblockDTMF = unblockDTMF;
this.playDTMF = playDTMF;
this.subscribeDTMF = subscribeDTMF;
this.unsubscribeDTMF = unsubscribeDTMF;
this.subscribeRequest = subscribeRequest;
this.subscribeAnswer = subscribeAnswer;
this.unsubscribe = unsubscribe;
const featureServer = await this.getFeatureServer();
if (!featureServer) {
@@ -63,7 +119,7 @@ class CallSession extends Emitter {
}
this.logger.debug(`using feature server ${featureServer}`);
this.rtpEngineOpts = makeRtpEngineOpts(this.req, SdpWantsSrtp(this.req.body), false, this.isFromMSTeams);
this.rtpEngineOpts = makeRtpEngineOpts(this.req, SdpWantsSrtp(sdp), false, this.isFromMSTeams);
this.rtpEngineResource = {destroy: this.del.bind(null, this.rtpEngineOpts.common)};
const obj = parseUri(this.req.uri);
let proxy, host, uri;
@@ -88,7 +144,7 @@ class CallSession extends Emitter {
...this.rtpEngineOpts.uac.mediaOpts,
'from-tag': this.rtpEngineOpts.uas.tag,
direction: ['public', 'private'],
sdp: this.req.body
sdp
};
const response = await this.offer(opts);
this.logger.debug({opts, response}, 'response from rtpengine to offer');
@@ -97,16 +153,26 @@ class CallSession extends Emitter {
throw new Error('rtpengine failed: answer');
}
// now send the INVITE in towards the feature servers
const headers = {
let headers = {
'From': createBLegFromHeader(this.req),
'To': this.req.get('To'),
'X-Account-Sid': this.req.locals.account_sid,
'X-CID': this.req.get('Call-ID'),
'X-Forwarded-For': `${this.req.source_address}:${this.req.source_port}`
'X-Forwarded-For': `${this.req.source_address}`
};
if (this.privateSipAddress) headers = {...headers, Contact: `<sip:${this.privateSipAddress}>`};
const spdOfferB = this.siprec && this.xml ?
createSiprecBody(headers, response.sdp, this.xml.type, this.xml.content) :
response.sdp;
const responseHeaders = {};
if (this.req.locals.carrier) Object.assign(headers, {'X-Originating-Carrier': this.req.locals.carrier});
if (this.req.locals.carrier) {
Object.assign(headers, {
'X-Originating-Carrier': this.req.locals.carrier,
'X-Voip-Carrier-Sid': this.req.locals.voip_carrier_sid
});
}
if (this.req.locals.msTeamsTenantFqdn) {
Object.assign(headers, {'X-MS-Teams-Tenant-FQDN': this.req.locals.msTeamsTenantFqdn});
@@ -131,14 +197,22 @@ class CallSession extends Emitter {
if (this.req.canceled) throw new Error('call canceled');
// now send the INVITE in towards the feature servers
debug(`sending INVITE to ${proxy} with ${uri}`);
const {uas, uac} = await this.srf.createB2BUA(this.req, this.res, uri, {
proxy,
headers,
responseHeaders,
proxyRequestHeaders: ['all', '-Authorization', '-Max-Forwards', '-Record-Route', '-Session-Expires', 'Min-SE'],
proxyResponseHeaders: ['all'],
localSdpB: response.sdp,
proxyRequestHeaders: [
'all',
'-Authorization',
'-Max-Forwards',
'-Record-Route',
'-Session-Expires',
'-X-Subspace-Forwarded-For'
],
proxyResponseHeaders: ['all', '-X-Trace-ID'],
localSdpB: spdOfferB,
localSdpA: async(sdp, res) => {
this.rtpEngineOpts.uac.tag = res.getParsedHeader('To').params.tag;
const opts = {
@@ -163,29 +237,40 @@ class CallSession extends Emitter {
this._setHandlers({uas, uac});
return;
} catch (err) {
this.rtpEngineResource.destroy();
this.rtpEngineResource.destroy().catch((err) => this.logger.info({err}, 'Error destroying rtpe after failure'));
this.activeCallIds.delete(this.req.get('Call-ID'));
this.stats.gauge('sbc.sip.calls.count', this.activeCallIds.size);
if (err instanceof SipError) {
const tags = ['accepted:no', `sipStatus:${err.status}`, `originator:${this.req.locals.originator}`];
this.stats.increment('sbc.terminations', tags);
this.logger.info(`call failed to connect to feature server with ${err.status}`);
return this.emit('failed');
this.emit('failed');
}
else if (err.message !== 'call canceled') {
this.logger.error(err, 'unexpected error routing inbound call');
}
this.srf.endSession(this.req);
}
}
_setDlgHandlers(dlg) {
this.activeCallIds.set(this.req.get('Call-ID'), this);
const {callId} = dlg.sip;
this.activeCallIds.set(callId, this);
this.subscribeDTMF(this.logger, callId, this.rtpEngineOpts.uas.tag,
this._onDTMF.bind(this));
dlg.on('destroy', () => {
debug('call ended with normal termination');
this.logger.info('call ended with normal termination');
this.rtpEngineResource.destroy().catch((err) => {});
this.activeCallIds.delete(this.req.get('Call-ID'));
this.activeCallIds.delete(callId);
if (dlg.other && dlg.other.connected) dlg.other.destroy().catch((e) => {});
if (this.srsClient) {
this.srsClient.stop();
this.srsClient = null;
}
this.srf.endSession(this.req);
});
//re-invite
@@ -202,22 +287,30 @@ class CallSession extends Emitter {
this.req.locals.cdr = {
...this.req.locals.cdr,
answered: true,
answered_at: callStart
answered_at: callStart,
trace_id: uac.res?.get('X-Trace-ID') || '00000000000000000000000000000000'
};
}
this.uas = uas;
this.uac = uac;
[uas, uac].forEach((dlg) => {
dlg.on('destroy', () => {
this.logger.info('call ended with normal termination');
dlg.on('destroy', async() => {
const other = dlg.other;
this.rtpEngineResource.destroy().catch((err) => {});
this.activeCallIds.delete(this.req.get('Call-ID'));
dlg.other.destroy().catch((e) => {});
if (process.env.JAMBONES_HOSTING) {
try {
await other.destroy();
} catch (err) {}
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
if (process.env.JAMBONES_HOSTING || process.env.JAMBONES_TRACK_ACCOUNT_CALLS) {
const {account_sid} = this.req.locals;
this.decrKey(this.callCountKey)
.then((count) => this.logger.debug({key: this.callCountKey},
`after hangup there are ${count} active calls for this account`))
.then((count) => {
this.logger.info(
{key: this.callCountKey},
`after hangup there are ${count} active calls for this account`);
return this.req.srf.locals.writeCallCount({account_sid, calls_in_progress: count});
})
.catch((err) => this.logger.error({err}, 'Error decrementing call count'));
}
@@ -235,16 +328,64 @@ class CallSession extends Emitter {
trunk
}).catch((err) => this.logger.error({err}, 'Error writing cdr for completed call'));
}
/* de-link the 2 Dialogs for GC */
dlg.removeAllListeners();
other.removeAllListeners();
dlg.other = null;
other.other = null;
if (this.srsClient) {
this.srsClient.stop();
this.srsClient = null;
}
this.logger.info(`call ended with normal termination, there are ${this.activeCallIds.size} active`);
this.srf.endSession(this.req);
});
});
this.subscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag,
this._onDTMF.bind(this, uac));
uas.on('modify', this._onReinvite.bind(this, uas));
uac.on('modify', this._onReinvite.bind(this, uac));
uac.on('refer', this._onFeatureServerTransfer.bind(this, uac));
uas.on('refer', this._onRefer.bind(this, uas));
uas.on('info', this._onInfo.bind(this, uas));
uac.on('info', this._onInfo.bind(this, uac));
// default forwarding of other request types
forwardInDialogRequests(uas, ['info', 'notify', 'options', 'message']);
forwardInDialogRequests(uas, ['notify', 'options', 'message']);
}
async _onDTMF(dlg, payload) {
this.logger.info({payload}, '_onDTMF');
try {
let dtmf;
switch (payload.event) {
case 10:
dtmf = '*';
break;
case 11:
dtmf = '#';
break;
default:
dtmf = '' + payload.event;
break;
}
await dlg.request({
method: 'INFO',
headers: {
'Content-Type': 'application/dtmf-relay'
},
body: `Signal=${dtmf}
Duration=${payload.duration} `
});
} catch (err) {
this.logger.info({err}, 'Error sending INFO application/dtmf-relay');
}
}
/**
@@ -254,7 +395,19 @@ class CallSession extends Emitter {
*/
async replaces(req, res) {
try {
let opts = Object.assign(this.rtpEngineOpts.offer, {sdp: req.body});
const fromTag = this.rtpEngineOpts.uas.tag;
const toTag = this.rtpEngineOpts.uac.tag;
const offerMedia = this.rtpEngineOpts.uac.mediaOpts;
const answerMedia = this.rtpEngineOpts.uas.mediaOpts;
const direction = ['public', 'private'];
let opts = {
...this.rtpEngineOpts.common,
...offerMedia,
'from-tag': fromTag,
'to-tag': toTag,
direction,
sdp: req.body,
};
let response = await this.offer(opts);
if ('ok' !== response.result) {
res.send(488);
@@ -262,8 +415,13 @@ class CallSession extends Emitter {
}
this.logger.info({opts, response}, 'sent offer for reinvite to rtpengine');
const sdp = await this.uac.modify(response.sdp);
opts = Object.assign(this.rtpEngineOpts.answer, {sdp, 'to-tag': this.toTag});
Object.assign(this.rtpEngineOpts.offer, {'to-tag': this.toTag});
opts = {
...this.rtpEngineOpts.common,
...answerMedia,
'from-tag': fromTag,
'to-tag': toTag,
sdp
};
response = await this.answer(opts);
if ('ok' !== response.result) {
res.send(488);
@@ -281,6 +439,8 @@ class CallSession extends Emitter {
});
}
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
const uas = await this.srf.createUAS(req, res, {
localSdp: response.sdp,
headers
@@ -301,26 +461,50 @@ class CallSession extends Emitter {
async _onReinvite(dlg, req, res) {
try {
/* check for re-invite with no SDP -- seen that from BT when they provide UUI info */
if (!req.body) {
this.logger.info('got a reINVITE with no SDP; just respond with our current offer');
res.send(200, {body: dlg.local.sdp});
return;
}
const 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 fromTag = dlg.type === 'uas' ? this.rtpEngineOpts.uas.tag : this.rtpEngineOpts.uac.tag;
const toTag = dlg.type === 'uas' ? this.rtpEngineOpts.uac.tag : this.rtpEngineOpts.uas.tag;
const offerMedia = dlg.type === 'uas' ? this.rtpEngineOpts.uac.mediaOpts : this.rtpEngineOpts.uas.mediaOpts;
const answerMedia = dlg.type === 'uas' ? this.rtpEngineOpts.uas.mediaOpts : this.rtpEngineOpts.uac.mediaOpts;
const direction = dlg.type === 'uas' ? ['private', 'public'] : ['public', 'private'];
const direction = dlg.type === 'uas' ? ['public', 'private'] : ['private', 'public'];
let opts = {
...this.rtpEngineOpts.common,
...offerMedia,
'from-tag': fromTag,
'to-tag': toTag,
direction,
sdp: req.body,
sdp: offeredSdp,
};
if (reason && opts.flags && !opts.flags.includes('reset')) opts.flags.push('reset');
let response = await this.offer(opts);
if ('ok' !== response.result) {
res.send(488);
throw new Error(`_onReinvite: rtpengine failed: offer: ${JSON.stringify(response)}`);
}
const sdp = await dlg.other.modify(response.sdp);
/* if this is a re-invite from the FS to change media anchoring, avoid sending the reinvite out */
let sdp;
if (reason && dlg.type === 'uac' && ['release-media', 'anchor-media'].includes(reason) &&
!this.callerIsUsingSrtp) {
this.logger.info({response}, `got a reinvite from FS to ${reason}`);
sdp = dlg.other.remote.sdp;
answerMedia.flags = ['asymmetric', 'port latching'];
this._mediaReleased = 'release-media' === reason;
}
else {
sdp = await dlg.other.modify(response.sdp);
}
opts = {
...this.rtpEngineOpts.common,
...answerMedia,
@@ -339,15 +523,219 @@ class CallSession extends Emitter {
}
}
async _onInfo(dlg, req, res) {
const fromTag = dlg.type === 'uas' ? this.rtpEngineOpts.uas.tag : this.rtpEngineOpts.uac.tag;
const toTag = dlg.type === 'uas' ? this.rtpEngineOpts.uac.tag : this.rtpEngineOpts.uas.tag;
const contentType = req.get('Content-Type');
try {
if (dlg.type === 'uac' && req.has('X-Reason')) {
const reason = req.get('X-Reason');
const opts = {
...this.rtpEngineOpts.common,
flags: ['reset'],
'from-tag': fromTag
};
this.logger.info(`_onInfo: got request ${reason}`);
if (reason.startsWith('mute')) {
const response = Promise.all([this.blockMedia(opts), this.blockDTMF(opts)]);
res.send(200);
this.logger.info({response}, `_onInfo: response to rtpengine command for ${reason}`);
}
else if (reason.startsWith('unmute')) {
const response = Promise.all([this.unblockMedia(opts), this.unblockDTMF(opts)]);
res.send(200);
this.logger.info({response}, `_onInfo: response to rtpengine command for ${reason}`);
}
else if (reason.includes('CallRecording')) {
let succeeded = false;
if (reason === 'startCallRecording') {
const from = this.req.getParsedHeader('From');
const to = this.req.getParsedHeader('To');
const aorFrom = from.uri;
const aorTo = to.uri;
this.logger.info({to, from}, 'startCallRecording request for a call');
const srsUrl = req.get('X-Srs-Url');
const srsRecordingId = req.get('X-Srs-Recording-ID');
const callSid = req.get('X-Call-Sid');
const accountSid = req.get('X-Account-Sid');
const applicationSid = req.get('X-Application-Sid');
if (this.srsClient) {
res.send(400);
this.logger.info('discarding duplicate startCallRecording request for a call');
return;
}
if (!srsUrl) {
this.logger.info('startCallRecording request is missing X-Srs-Url header');
res.send(400);
return;
}
this.srsClient = new SrsClient(this.logger, {
srf: dlg.srf,
direction: 'inbound',
originalInvite: this.req,
callingNumber: this.req.callingNumber,
calledNumber: this.req.calledNumber,
srsUrl,
srsRecordingId,
callSid,
accountSid,
applicationSid,
rtpEngineOpts: this.rtpEngineOpts,
fromTag,
toTag,
aorFrom,
aorTo,
subscribeRequest: this.subscribeRequest,
subscribeAnswer: this.subscribeAnswer,
del: this.del,
blockMedia: this.blockMedia,
unblockMedia: this.unblockMedia,
unsubscribe: this.unsubscribe
});
try {
succeeded = await this.srsClient.start();
} catch (err) {
this.logger.error({err}, 'Error starting SipRec call recording');
}
}
else if (reason === 'stopCallRecording') {
if (!this.srsClient) {
res.send(400);
this.logger.info('discarding stopCallRecording request because we are not recording');
return;
}
try {
succeeded = await this.srsClient.stop();
} catch (err) {
this.logger.error({err}, 'Error stopping SipRec call recording');
}
this.srsClient = null;
}
else if (reason === 'pauseCallRecording') {
if (!this.srsClient || this.srsClient.paused) {
this.logger.info('discarding invalid pauseCallRecording request');
res.send(400);
return;
}
succeeded = await this.srsClient.pause();
}
else if (reason === 'resumeCallRecording') {
if (!this.srsClient || !this.srsClient.paused) {
res.send(400);
this.logger.info('discarding invalid resumeCallRecording request');
return;
}
succeeded = await this.srsClient.resume();
}
res.send(succeeded ? 200 : 503);
}
}
else if (dlg.type === 'uas' && ['application/dtmf-relay', 'application/dtmf'].includes(contentType)) {
const arr = /Signal=\s*([1-9#*])/.exec(req.body);
if (!arr) {
this.logger.info({body: req.body}, '_onInfo: invalid INFO dtmf request');
throw new Error(`_onInfo: no dtmf in body for ${contentType}`);
}
const code = arr[1];
const arr2 = /Duration=\s*(\d+)/.exec(req.body);
const duration = arr2 ? arr2[1] : 250;
if (this.isMediaReleased) {
/* just relay on to the feature server */
this.logger.info({code, duration}, 'got SIP INFO DTMF from caller, relaying to feature server');
this._onDTMF(dlg.other, {event: code, duration})
.catch((err) => this.logger.info({err}, 'Error relaying DTMF to feature server'));
res.send(200);
}
else {
/* else convert SIP INFO to RFC 2833 telephony events */
this.logger.info({code, duration}, 'got SIP INFO DTMF from caller, converting to RFC 2833');
const opts = {
...this.rtpEngineOpts.common,
'from-tag': this.rtpEngineOpts.uas.tag,
code,
duration
};
const response = await this.playDTMF(opts);
if ('ok' !== response.result) {
this.logger.info({response}, `rtpengine playDTMF failed with ${JSON.stringify(response)}`);
throw new Error('rtpengine failed: answer');
}
res.send(200);
}
}
else {
const immutableHdrs = ['via', 'from', 'to', 'call-id', 'cseq', 'max-forwards', 'content-length'];
const headers = {};
Object.keys(req.headers).forEach((h) => {
if (!immutableHdrs.includes(h)) headers[h] = req.headers[h];
});
const response = await dlg.other.request({method: 'INFO', headers, body: req.body});
const responseHeaders = {};
if (response.has('Content-Type')) {
Object.assign(responseHeaders, {'Content-Type': response.get('Content-Type')});
}
res.send(response.status, {headers: responseHeaders, body: response.body});
}
} catch (err) {
if (this.srsClient) {
this.srsClient = null;
}
res.send(500);
this.logger.info({err}, `Error handing INFO request on ${dlg.type} leg`);
}
}
async _onFeatureServerTransfer(dlg, req, res) {
try {
const referTo = req.getParsedHeader('Refer-To');
const uri = parseUri(referTo.uri);
this.logger.info({uri, referTo}, 'received REFER from feature server');
this.logger.info({uri, referTo, headers: req.headers}, 'received REFER from feature server');
const arr = /context-(.*)/.exec(uri.user);
if (!arr) {
this.logger.info(`invalid Refer-To header: ${referTo.uri}`);
return res.send(501);
/* call transfer requested */
const {gateway} = this.req.locals;
const referredBy = req.getParsedHeader('Referred-By');
if (!referredBy) return res.send(400);
const u = parseUri(referredBy.uri);
let selectedGateway = false;
let e164 = false;
if (gateway) {
/* host of Refer-to to an outbound gateway */
const gw = await this.srf.locals.getOutboundGatewayForRefer(gateway.voip_carrier_sid);
if (gw) {
selectedGateway = true;
e164 = gw.e164_leading_plus;
uri.host = gw.ipv4;
uri.port = gw.port;
}
}
if (!selectedGateway) {
uri.host = this.req.source_address;
uri.port = this.req.source_port;
}
if (e164 && !uri.user.startsWith('+')) {
uri.user = `+${uri.user}`;
}
// eslint-disable-next-line no-unused-vars
const {via, from, to, 'call-id':callid, cseq, 'max-forwards':maxforwards,
// eslint-disable-next-line no-unused-vars
'content-length':contentlength, 'refer-to':_referto, 'referred-by':_referredby,
...customHeaders
} = req.headers;
const response = await this.uas.request({
method: 'REFER',
headers: {
'Refer-To': stringifyUri(uri),
'Referred-By': stringifyUri(u),
...customHeaders
}
});
return res.send(response.status);
}
res.send(202);
@@ -367,6 +755,7 @@ class CallSession extends Emitter {
this.rtpEngineResource.destroy();
this.activeCallIds.delete(this.req.get('Call-ID'));
uac.other.destroy();
this.srf.endSession(this.req);
});
// now we can destroy the old dialog
dlg.destroy().catch(() => {});
@@ -438,6 +827,7 @@ class CallSession extends Emitter {
// successfully connected
this.logger.info('successfully connected new call leg for REFER');
this.unsubscribeDTMF(this.logger, this.req.get('Call-ID'), this.rtpEngineOpts.uas.tag);
this.referInvite = null;
sendNotify(this.uas, '200 OK');
this.uas.destroy();
@@ -451,8 +841,20 @@ class CallSession extends Emitter {
}
}
else {
// TODO: forward on to feature server
res.send(501);
/* REFER coming in from a sip device, forward to feature server */
try {
const response = await dlg.other.request({
method: 'REFER',
headers: {
'Refer-To': req.get('Refer-To'),
'Referred-By': req.get('Referred-By'),
'User-Agent': req.get('User-Agent')
}
});
res.send(response.status, response.reason);
} catch (err) {
this.logger.error({err}, 'CallSession:_onRefer: error handling incoming REFER');
}
}
}
+7
View File
@@ -0,0 +1,7 @@
{
"LifeCycleEvents" : {
"ScaleIn": "scale-in",
"StandbyEnter": "standby-enter",
"StandbyExit": "standby-exit"
}
}
+71 -7
View File
@@ -11,6 +11,15 @@ WHERE acc.sip_realm = ?
AND vc.account_sid = acc.account_sid
AND sg.voip_carrier_sid = vc.voip_carrier_sid`;
const sqlSelectAllCarriersForSPByRealm =
`SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.account_sid,
vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask
FROM sip_gateways sg, voip_carriers vc, accounts acc
WHERE acc.sip_realm = ?
AND vc.service_provider_sid = acc.service_provider_sid
AND vc.account_sid IS NULL
AND sg.voip_carrier_sid = vc.voip_carrier_sid`;
const sqlSelectAllGatewaysForSP =
`SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.service_provider_sid,
vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask
@@ -35,6 +44,17 @@ SELECT * FROM phone_numbers
WHERE number = ?
AND voip_carrier_sid = ?`;
const sqlQueryAllDidsForCarrier = `
SELECT * FROM phone_numbers
WHERE voip_carrier_sid = ?`;
const sqlSelectOutboundGatewayForCarrier = `
SELECT ipv4, port, e164_leading_plus
FROM sip_gateways sg, voip_carriers vc
WHERE sg.voip_carrier_sid = ?
AND sg.voip_carrier_sid = vc.voip_carrier_sid
AND outbound = 1`;
const gatewayMatchesSourceAddress = (source_address, gw) => {
if (32 === gw.netmask && gw.ipv4 === source_address) return true;
if (gw.netmask < 32) {
@@ -48,13 +68,35 @@ module.exports = (srf, logger) => {
const {pool} = srf.locals.dbHelpers;
const pp = pool.promise();
const getOutboundGatewayForRefer = async(voip_carrier_sid) => {
try {
const [r] = await pp.query(sqlSelectOutboundGatewayForCarrier, [voip_carrier_sid]);
if (0 === r.length) return null;
/* if multiple, prefer a DNS name */
const hasDns = r.find((row) => row.ipv4.match(/^[A-Za-z]/));
return hasDns || r[0];
} catch (err) {
logger.error({err}, 'getOutboundGatewayForRefer');
}
};
const getApplicationForDidAndCarrier = async(req, voip_carrier_sid) => {
const did = normalizeDID(req.calledNumber);
try {
/* straight DID match */
const [r] = await pp.query(sqlQueryApplicationByDid, [did, voip_carrier_sid]);
if (0 === r.length) return null;
return r[0].application_sid;
if (r.length) return r[0].application_sid;
/* wildcard / regex match */
const [r2] = await pp.query(sqlQueryAllDidsForCarrier, [voip_carrier_sid]);
const match = r2
.filter((o) => o.number.match(/\D/)) // look at anything with non-digit characters
.sort((a, b) => b.number.length - a.number.length) // prefer longest match
.find((o) => did.match(new RegExp(o.number.endsWith('*') ? `${o.number.slice(0, -1)}\\d*` : o.number)));
if (match) return match.application_sid;
return null;
} catch (err) {
logger.error({err}, 'getApplicationForDidAndCarrier');
}
@@ -73,7 +115,11 @@ module.exports = (srf, logger) => {
const [r] = await pp.query(sqlCarriersForAccountBySid,
[process.env.SBC_ACCOUNT_SID, req.source_address, req.source_port]);
if (0 === r.length) return failure;
return {fromCarrier: true, gateway: r[0], account_sid: process.env.SBC_ACCOUNT_SID};
return {
fromCarrier: true,
gateway: r[0],
account_sid: process.env.SBC_ACCOUNT_SID
};
}
else {
/* we may have a carrier at the service provider level */
@@ -88,8 +134,23 @@ module.exports = (srf, logger) => {
const [r] = await pp.query(sql);
if (0 === r.length) {
/* came from a carrier, but number is not provisioned */
return {fromCarrier: true};
/* came from a carrier, but number is not provisioned..
check if we only have a single account, otherwise we have no
way of knowing which account this is for
*/
const [r] = await pp.query('SELECT count(*) as count from accounts where service_provider_sid = ?',
matches[0].service_provider_sid);
if (r[0].count === 0) return {fromCarrier: true};
else {
const [accounts] = await pp.query('SELECT * from accounts where service_provider_sid = ?',
matches[0].service_provider_sid);
return {
fromCarrier: true,
gateway: matches[0],
account_sid: accounts[0].account_sid,
account: accounts[0]
};
}
}
const gateway = matches.find((m) => m.voip_carrier_sid === r[0].voip_carrier_sid);
const [accounts] = await pp.query(sqlAccountBySid, r[0].account_sid);
@@ -107,7 +168,9 @@ module.exports = (srf, logger) => {
}
/* get all the carriers and gateways for the account owning this sip realm */
const [gw] = await pp.query(sqlSelectAllCarriersForAccountByRealm, uri.host);
const [gwAcc] = await pp.query(sqlSelectAllCarriersForAccountByRealm, uri.host);
const [gwSP] = gwAcc.length ? [[]] : await pp.query(sqlSelectAllCarriersForSPByRealm, uri.host);
const gw = gwAcc.concat(gwSP);
const selected = gw.find(gatewayMatchesSourceAddress.bind(null, req.source_address));
if (selected) {
const [a] = await pp.query(sqlAccountByRealm, uri.host);
@@ -125,6 +188,7 @@ module.exports = (srf, logger) => {
return {
wasOriginatedFromCarrier,
getApplicationForDidAndCarrier
getApplicationForDidAndCarrier,
getOutboundGatewayForRefer
};
};
+15 -6
View File
@@ -1,4 +1,8 @@
const setName = `${(process.env.JAMBONES_CLUSTER_ID || 'default')}:active-fs`;
const assert = require('assert');
assert.ok(!process.env.K8S || process.env.K8S_FEATURE_SERVER_SERVICE_NAME,
'when running in Kubernetes, an env var K8S_FEATURE_SERVER_SERVICE_NAME is required');
module.exports = (srf, logger) => {
const {retrieveSet, createSet} = srf.locals.realtimeDbHelpers;
@@ -10,13 +14,18 @@ module.exports = (srf, logger) => {
return async() => {
try {
const fs = await retrieveSet(setName);
if (0 === fs.length) {
logger.info('No available feature servers to handle incoming call');
return;
if (process.env.K8S) {
return process.env.K8S_FEATURE_SERVER_SERVICE_NAME;
}
else {
const fs = await retrieveSet(setName);
if (0 === fs.length) {
logger.info('No available feature servers to handle incoming call');
return;
}
logger.debug({fs}, `retrieved ${setName}`);
return fs[idx++ % fs.length];
}
logger.debug({fs}, `retrieved ${setName}`);
return fs[idx++ % fs.length];
} catch (err) {
logger.error({err}, `Error retrieving ${setName}`);
}
+78 -16
View File
@@ -68,13 +68,22 @@ module.exports = function(srf, logger) {
blacklistUnknownRealms: true,
emitter: new AuthOutcomeReporter(stats)
});
const {wasOriginatedFromCarrier, getApplicationForDidAndCarrier} = require('./db-utils')(srf, logger);
const initLocals = (req, res, next) => {
req.locals = req.locals || {};
req.locals.cdr = initCdr(req);
const callId = req.get('Call-ID');
req.locals = req.locals || {callId};
/* check if forwarded by a proxy that applied an X-Forwarded-For Header */
if (req.has('X-Forwarded-For') || req.has('X-Subspace-Forwarded-For')) {
const original_source_address = req.get('X-Forwarded-For') || req.get('X-Subspace-Forwarded-For');
logger.info({
callId: req.get('Call-ID'),
original_source_address,
proxy_source_address: req.source_address,
}, 'overwriting source address for proxied SIP INVITE');
req.source_address = original_source_address;
}
req.locals.cdr = initCdr(req);
req.on('cancel', () => {
logger.info({callId}, 'caller hungup before connecting to feature server');
req.canceled = true;
@@ -100,10 +109,39 @@ module.exports = function(srf, logger) {
next();
};
const handleSipRec = async(req, res, next) => {
const {callId} = req.locals;
if (Array.isArray(req.payload) && req.payload.length > 1) {
const sdp = req.payload
.find((p) => p.type === 'application/sdp')
.content;
if (!sdp) {
logger.error({callId}, 'No SDP in multipart sdp');
return res.send(503);
}
const xml = req.payload.find((p) => p.type !== 'application/sdp');
const endPos = xml.content.indexOf('</recording>');
xml.content = endPos !== -1 ?
`${xml.content.substring(0, endPos + 12)}` :
xml.content;
logger.debug({callId, xml}, 'incoming call with SIPREC body');
req.locals = {...req.locals, sdp, siprec: true, xml};
}
else req.locals = {...req.locals, sdp: req.body};
next();
};
const identifyAccount = async(req, res, next) => {
try {
const {fromCarrier, gateway, account_sid, application_sid, account} = await wasOriginatedFromCarrier(req);
const {siprec, callId} = req.locals;
const {wasOriginatedFromCarrier, getApplicationForDidAndCarrier} = req.srf.locals;
const {
fromCarrier,
gateway,
account_sid,
application_sid,
account
} = await wasOriginatedFromCarrier(req);
/**
* calls come from 3 sources:
* (1) A carrier
@@ -117,11 +155,23 @@ module.exports = function(srf, logger) {
}
logger.debug({gateway}, 'identifyAccount: incoming call from gateway');
/* check for phone number level routing */
const sid = application_sid || await getApplicationForDidAndCarrier(req, gateway.voip_carrier_sid);
let sid;
if (siprec) {
if (!account.siprec_hook_sid) {
logger.info({callId}, 'identifyAccount: rejecting call because SIPREC hook has not been provisioned');
return res.send(404);
}
sid = account.siprec_hook_sid;
}
else {
/* check for phone number level routing */
sid = application_sid || await getApplicationForDidAndCarrier(req, gateway.voip_carrier_sid);
}
req.locals = {
originator: 'trunk',
carrier: gateway.name,
gateway,
voip_carrier_sid: gateway.voip_carrier_sid,
application_sid: sid || gateway.application_sid,
account_sid,
account,
@@ -135,7 +185,8 @@ module.exports = function(srf, logger) {
const app = await lookupAppByTeamsTenant(uri.host);
if (!app) {
stats.increment('sbc.terminations', ['sipStatus:404']);
return res.send(404, {headers: {'X-Reason': 'no configured application'}});
res.send(404, {headers: {'X-Reason': 'no configured application'}});
return req.srf.endSession(req);
}
req.locals = {
@@ -154,7 +205,8 @@ module.exports = function(srf, logger) {
const account = await lookupAccountBySipRealm(uri.host);
if (!account) {
stats.increment('sbc.terminations', ['sipStatus:404']);
return res.send(404);
res.send(404);
return req.srf.endSession(req);
}
/* if this is a dedicated SBC (static IP) only take calls for that account's sip realm */
@@ -163,7 +215,8 @@ module.exports = function(srf, logger) {
`identifyAccount: static IP for ${process.env.SBC_ACCOUNT_SID} but call for ${account.account_sid}`);
stats.increment('sbc.terminations', ['sipStatus:404']);
delete req.locals.cdr;
return res.send(404);
res.send(404);
return req.srf.endSession(req);
}
req.locals = {
account_sid: account.account_sid,
@@ -200,11 +253,11 @@ module.exports = function(srf, logger) {
};
const checkLimits = async(req, res, next) => {
if (!process.env.JAMBONES_HOSTING) return next(); // skip
if (!process.env.JAMBONES_HOSTING && !process.env.JAMBONES_TRACK_ACCOUNT_CALLS) return next(); // skip
const {incrKey, decrKey} = req.srf.locals.realtimeDbHelpers;
const {logger, account_sid} = req.locals;
const {writeAlerts, AlertType} = req.srf.locals;
const {writeCallCount, writeAlerts, AlertType} = req.srf.locals;
assert(account_sid);
const key = makeCallCountKey(account_sid);
@@ -213,10 +266,11 @@ module.exports = function(srf, logger) {
if (status > 200) {
decrKey(key)
.then((count) => {
logger.debug({key}, `after rejection there are ${count} active calls for this account`);
logger.info({key}, `after rejection there are ${count} active calls for this account`);
debug({key}, `after rejection there are ${count} active calls for this account`);
return;
return count;
})
.then((count) => writeCallCount({account_sid, calls_in_progress: count}))
.catch((err) => logger.error({err}, 'checkLimits: decrKey err'));
}
});
@@ -224,6 +278,10 @@ module.exports = function(srf, logger) {
try {
/* increment the call count */
const calls = await incrKey(key);
writeCallCount({account_sid, calls_in_progress: calls})
.then(() => logger.info(`checkLimits: after incrementing there are ${calls} active calls for this account`))
.catch((err) => logger.error({err}, 'checkLimits: error writing call count'));
if (!process.env.JAMBONES_HOSTING) return next();
/* compare to account's limit, though avoid db hit when call count is low */
const minLimit = process.env.MIN_CALL_LIMIT ?
@@ -244,13 +302,15 @@ module.exports = function(srf, logger) {
account_sid,
count: limit_sessions
}).catch((err) => logger.info({err}, 'checkLimits: error writing alert'));
return res.send(503, 'Maximum Calls In Progress');
res.send(503, 'Maximum Calls In Progress');
return req.srf.endSession(req);
}
next();
} catch (err) {
stats.increment('sbc.terminations', ['sipStatus:500']);
logger.error({err}, 'error checking limits error for inbound call');
res.send(500);
req.srf.endSession(req);
}
};
@@ -263,11 +323,13 @@ module.exports = function(srf, logger) {
stats.increment('sbc.terminations', ['sipStatus:500']);
logger.error(err, `${req.get('Call-ID')} Error looking up related info for inbound call`);
res.send(500);
req.srf.endSession(req);
}
};
return {
initLocals,
handleSipRec,
challengeDeviceCalls,
identifyAccount,
checkLimits
+43 -3
View File
@@ -15,7 +15,13 @@ 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 srctOpts = srcIsUsingSrtp ? srtpOpts : rtpCharacteristics;
const srcOpts = srcIsUsingSrtp ? srtpOpts : rtpCharacteristics;
/* 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']
@@ -24,7 +30,7 @@ function makeRtpEngineOpts(req, srcIsUsingSrtp, dstIsUsingSrtp, teams = false) {
common,
uas: {
tag: from.params.tag,
mediaOpts: srctOpts
mediaOpts: srcOpts
},
uac: {
tag: null,
@@ -45,11 +51,45 @@ const normalizeDID = (tel) => {
return arr ? arr[1] : tel;
};
const equalsIgnoreOrder = (a, b) => {
if (a.length !== b.length) return false;
const uniqueValues = new Set([...a, ...b]);
for (const v of uniqueValues) {
const aCount = a.filter((e) => e === v).length;
const bCount = b.filter((e) => e === v).length;
if (aCount !== bCount) return false;
}
return true;
};
const systemHealth = async(redisClient, ping, getCount) => {
await Promise.all([redisClient.ping(), ping()]);
return getCount();
};
const 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);
});
});
};
module.exports = {
isWSS,
SdpWantsSrtp,
getAppserver,
makeRtpEngineOpts,
makeCallCountKey,
normalizeDID
normalizeDID,
equalsIgnoreOrder,
systemHealth,
createHealthCheckApp
};
+24
View File
@@ -0,0 +1,24 @@
#!/bin/sh
TCP_SERVER_PORT="${DRACHTIO_PORT:-4000}"
nc -v -z localhost $TCP_SERVER_PORT
# if last command exited with non zero
if [ $? != 0 ]
then
exit 1
fi
HTTP_SERVER_PORT="${HTTP_PORT:-3000}"
printf 'GET /system-health HTTP/1.1\r\nHost: localhost\r\n\r\n' | nc -v localhost 3000 | grep calls
# grep will automatically exit with 1 if string is not matched, however, will leave that call there in case
# we pivot to pipe to dev/null
if [ $? != 0 ]
then
exit 1
fi
exit 0
+3173 -1127
View File
File diff suppressed because it is too large Load Diff
+24 -17
View File
@@ -1,9 +1,9 @@
{
"name": "sbc-inbound",
"version": "0.3.6",
"version": "v0.7.6",
"main": "app.js",
"engines": {
"node": ">= 10.16.0"
"node": ">= 12.0.0"
},
"keywords": [
"sip",
@@ -20,29 +20,36 @@
},
"scripts": {
"start": "node app",
"test": "NODE_ENV=test JAMBONES_HOSTING=1 SBC_ACCOUNT_SID=ed649e33-e771-403a-8c99-1780eabbc803 JAMBONES_TIME_SERIES_HOST=127.0.0.1 JAMBONES_MYSQL_HOST=127.0.0.1 JAMBONES_MYSQL_USER=jambones_test JAMBONES_MYSQL_PASSWORD=jambones_test JAMBONES_MYSQL_DATABASE=jambones_test JAMBONES_REDIS_HOST=localhost JAMBONES_REDIS_PORT=16379 JAMBONES_LOGLEVEL=debug DRACHTIO_SECRET=cymru DRACHTIO_HOST=127.0.0.1 DRACHTIO_PORT=9060 JAMBONES_RTPENGINES=127.0.0.1:12222 JAMBONES_FEATURE_SERVERS=172.38.0.11 node test/ ",
"test": "NODE_ENV=test HTTP_PORT=3050 JAMBONES_NETWORK_CIDR='127.0.0.1/32' JAMBONES_HOSTING=1 SBC_ACCOUNT_SID=ed649e33-e771-403a-8c99-1780eabbc803 JAMBONES_TIME_SERIES_HOST=127.0.0.1 JAMBONES_MYSQL_HOST=127.0.0.1 JAMBONES_MYSQL_USER=jambones_test JAMBONES_MYSQL_PASSWORD=jambones_test JAMBONES_MYSQL_DATABASE=jambones_test JAMBONES_REDIS_HOST=localhost JAMBONES_REDIS_PORT=16379 JAMBONES_LOGLEVEL=info DRACHTIO_SECRET=cymru DRACHTIO_HOST=127.0.0.1 DRACHTIO_PORT=9060 JAMBONES_RTPENGINES=127.0.0.1:12222 JAMBONES_FEATURE_SERVERS=172.38.0.11 node test/ ",
"coverage": "./node_modules/.bin/nyc --reporter html --report-dir ./coverage npm run test",
"jslint": "eslint app.js lib"
},
"dependencies": {
"@jambonz/db-helpers": "^0.6.12",
"@jambonz/http-authenticator": "^0.2.0",
"@jambonz/realtimedb-helpers": "^0.4.3",
"@jambonz/rtpengine-utils": "^0.1.12",
"@jambonz/stats-collector": "^0.1.5",
"@jambonz/time-series": "^0.1.5",
"@jambonz/db-helpers": "^0.6.18",
"@jambonz/http-authenticator": "^0.2.1",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/realtimedb-helpers": "^0.4.29",
"@jambonz/rtpengine-utils": "^0.3.1",
"@jambonz/siprec-client-utils": "^0.1.4",
"@jambonz/stats-collector": "^0.1.6",
"@jambonz/time-series": "^0.1.12",
"aws-sdk": "^2.1152.0",
"bent": "^7.3.12",
"cidr-matcher": "^2.1.1",
"debug": "^4.3.1",
"debug": "^4.3.4",
"drachtio-fn-b2b-sugar": "0.0.12",
"drachtio-srf": "^4.4.49",
"pino": "^6.8.0",
"rtpengine-client": "^0.2.0"
"drachtio-srf": "^4.5.1",
"express": "^4.18.1",
"pino": "^7.11.0",
"sdp-transform": "^2.14.1",
"uuid": "^8.3.2",
"verify-aws-sns-signature": "^0.0.7",
"xml2js": "^0.4.23"
},
"devDependencies": {
"clear-module": "^4.1.1",
"eslint": "^7.15.0",
"eslint-plugin-promise": "^4.2.1",
"eslint": "^7.32.0",
"eslint-plugin-promise": "^4.3.1",
"nyc": "^15.1.0",
"tape": "^4.13.3"
"tape": "^4.15.1"
}
}
+7
View File
@@ -56,5 +56,12 @@ values ('888a5339-c62c-4075-9e19-f4de70a96597', '999c1452-620d-4195-9f19-c9814ef
insert into phone_numbers (phone_number_sid, number, voip_carrier_sid, account_sid)
values ('999a5339-c62c-4075-9e19-f4de70a96597', '16173333456', '287c1452-620d-4195-9f19-c9814ef90d78', 'ed649e33-e771-403a-8c99-1780eabbc803');
insert into phone_numbers (phone_number_sid, number, voip_carrier_sid, account_sid)
values ('29543d4e-d959-4a25-836a-cde7161cd7d5', '1508222*', '287c1452-620d-4195-9f19-c9814ef90d78', 'ed649e33-e771-403a-8c99-1780eabbc803');
insert into phone_numbers (phone_number_sid, number, voip_carrier_sid, account_sid)
values ('dddd5c34-feae-4d70-98af-bb4d1f8dc965', '1508*', '287c1452-620d-4195-9f19-c9814ef90d78', 'ed649e33-e771-403a-8c99-1780eabbc803');
insert into phone_numbers (phone_number_sid, number, voip_carrier_sid, account_sid)
values ('d458bf7a-bcea-47b2-ac96-66dfc9c5c220', '150822233*', '287c1452-620d-4195-9f19-c9814ef90d78', 'ed649e33-e771-403a-8c99-1780eabbc803');
insert into phone_numbers (phone_number_sid, number, voip_carrier_sid, account_sid)
values ('f7ad205d-b92f-4363-8160-f8b5216b40d3', '15083871234', '287c1452-620d-4195-9f19-c9814ef90d78', 'd7cc37cb-d152-49ef-a51b-485f6e917089');
+2 -1
View File
@@ -10,6 +10,7 @@ networks:
services:
mysql:
image: mysql:5.7
platform: linux/x86_64
ports:
- "3306:3306"
environment:
@@ -71,7 +72,7 @@ services:
ipv4_address: 172.38.0.14
influxdb:
image: influxdb:1.8-alpine
image: influxdb:1.8
ports:
- "8086:8086"
networks:
+118
View File
@@ -0,0 +1,118 @@
<?xml version="1.0" encoding="ISO-8859-1" ?>
<!DOCTYPE scenario SYSTEM "sipp.dtd">
<!-- This program is free software; you can redistribute it and/or -->
<!-- modify it under the terms of the GNU General Public License as -->
<!-- published by the Free Software Foundation; either version 2 of the -->
<!-- License, or (at your option) any later version. -->
<!-- -->
<!-- This program is distributed in the hope that it will be useful, -->
<!-- but WITHOUT ANY WARRANTY; without even the implied warranty of -->
<!-- MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the -->
<!-- GNU General Public License for more details. -->
<!-- -->
<!-- You should have received a copy of the GNU General Public License -->
<!-- along with this program; if not, write to the -->
<!-- Free Software Foundation, Inc., -->
<!-- 59 Temple Place, Suite 330, Boston, MA 02111-1307 USA -->
<!-- -->
<!-- Sipp 'uac' scenario with pcap (rtp) play -->
<!-- -->
<scenario name="UAC with media">
<!-- In client mode (sipp placing calls), the Call-ID MUST be -->
<!-- generated by sipp. To do so, use [call_id] keyword. -->
<send retrans="500">
<![CDATA[
INVITE sip:+15082223333@jambonz.org SIP/2.0
Via: SIP/2.0/[transport] [local_ip]:[local_port];branch=[branch]
From: sipp <sip:sipp@[local_ip]:[local_port]>;tag=[pid]SIPpTag09[call_number]
To: <sip:15082223333@jambonz.org>
Call-ID: [call_id]
CSeq: 1 INVITE
Contact: sip:sipp@[local_ip]:[local_port]
Max-Forwards: 70
Subject: uac-pcap-carrier-success
Content-Type: application/sdp
Content-Length: [len]
v=0
o=user1 53655765 2353687637 IN IP[local_ip_type] [local_ip]
s=-
c=IN IP[local_ip_type] [local_ip]
t=0 0
m=audio [auto_media_port] RTP/AVP 8 101
a=rtpmap:8 PCMA/8000
a=rtpmap:101 telephone-event/8000
a=fmtp:101 0-11,16
]]>
</send>
<recv response="100" optional="true">
</recv>
<recv response="180" optional="true">
</recv>
<!-- By adding rrs="true" (Record Route Sets), the route sets -->
<!-- are saved and used for following messages sent. Useful to test -->
<!-- against stateful SIP proxies/B2BUAs. -->
<recv response="200" rtd="true" crlf="true">
</recv>
<!-- Packet lost can be simulated in any send/recv message by -->
<!-- by adding the 'lost = "10"'. Value can be [1-100] percent. -->
<send>
<![CDATA[
ACK sip:15082223333@jambonz.org SIP/2.0
Via: SIP/2.0/[transport] [local_ip]:[local_port];branch=[branch]
From: sipp <sip:sipp@[local_ip]:[local_port]>;tag=[pid]SIPpTag09[call_number]
To: <sip:15082223333@jambonz.org>[peer_tag_param]
Call-ID: [call_id]
CSeq: 1 ACK
Max-Forwards: 70
Subject: uac-pcap-carrier-success
Content-Length: 0
]]>
</send>
<!-- Play a pre-recorded PCAP file (RTP stream) -->
<nop>
<action>
<exec play_pcap_audio="pcap/g711a.pcap"/>
</action>
</nop>
<!-- Pause briefly -->
<pause milliseconds="3000"/>
<!-- The 'crlf' option inserts a blank line in the statistics report. -->
<send retrans="500">
<![CDATA[
BYE sip:15082223333@jambonz.org SIP/2.0
Via: SIP/2.0/[transport] [local_ip]:[local_port];branch=[branch]
From: sipp <sip:sipp@[local_ip]:[local_port]>;tag=[pid]SIPpTag09[call_number]
To: <sip:15082223333@jambonz.org>[peer_tag_param]
Call-ID: [call_id]
CSeq: 2 BYE
Subject: uac-pcap-carrier-success
Content-Length: 0
]]>
</send>
<recv response="200" crlf="true">
</recv>
<!-- definition of the response time repartition table (unit is ms) -->
<ResponseTimeRepartition value="10, 20, 30, 40, 50, 100, 150, 200"/>
<!-- definition of the call length repartition table (unit is ms) -->
<CallLengthRepartition value="10, 50, 100, 500, 1000, 5000, 10000"/>
</scenario>
+18 -1
View File
@@ -124,7 +124,24 @@
]]>
</send>
<recv request="BYE">
</recv>
<send next="2">
<![CDATA[
SIP/2.0 200 OK
[last_Via:]
[last_From:]
[last_To:]
[last_Call-ID:]
[last_CSeq:]
Contact: <sip:[local_ip]:[local_port];transport=[transport]>
Content-Length: 0
]]>
</send>
<label id="2"/>
</scenario>
+13 -5
View File
@@ -1,8 +1,7 @@
const test = require('tape');
const { output, sippUac } = require('./sipp')('test_sbc-inbound');
const debug = require('debug')('drachtio:sbc-inbound');
const clearModule = require('clear-module');
const consoleLogger = {error: console.error, info: console.log, debug: console.log};
const { sippUac } = require('./sipp')('test_sbc-inbound');
const bent = require('bent');
const getJSON = bent('json');
process.on('unhandledRejection', (reason, p) => {
console.log('Unhandled Rejection at: Promise', p, 'reason:', reason);
@@ -28,12 +27,21 @@ test('incoming call 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)')
await sippUac('uac-pcap-carrier-success.xml', '172.38.0.20');
t.pass('incoming call from carrier completed successfully');
await sippUac('uac-pcap-pbx-success.xml', '172.38.0.21');
t.pass('incoming call from account-level carrier completed successfully');
await sippUac('uac-did-regex-match.xml', '172.38.0.20');
t.pass('incoming call matched by trailing wildcard *');
await sippUac('uac-pcap-device-success.xml', '172.38.0.30');
t.pass('incoming call from authenticated device completed successfully');
@@ -55,7 +63,7 @@ test('incoming call tests', async(t) => {
await waitFor(10);
const res = await queryCdrs({account_sid: 'ed649e33-e771-403a-8c99-1780eabbc803'});
console.log(`cdrs: ${JSON.stringify(res)}`);
t.ok(6 === res.total, 'successfully wrote 6 cdrs for calls');
t.ok(7 === res.total, 'successfully wrote 7 cdrs for calls');
srf.disconnect();
t.end();