Compare commits

...
29 Commits
Author SHA1 Message Date
Dave Horton d0d8ba93b2 0.9.11 2026-09-15 12:38:35 -04:00
Dave Horton 3f729e232e 0.9.10 2026-09-15 12:38:16 -04:00
Hoan Luu HuuandClaude Opus 5 509fdf6c37 Fix/siprec survives fs transfer (#249)
* fix: keep siprec recording alive across a feature server transfer

A cross-feature-server move (enqueue/dequeue or conference) re-negotiates the
feature-server leg in _onFeatureServerTransfer, but nothing rebuilt the rtpengine
subscription the SIPREC recording forks from, so the recorder went silent from the
moment the call moved. The fresh destroy handler installed on the new leg also
dropped the _stopRecording() call the original handlers have, so the SIPREC dialog
was never BYEd and the recorder had to wait out its media timeout.

Rebuild the subscription after the transfer re-negotiates media, and stop the
recording when the transferred leg ends. The resubscribe call is guarded so this
is safe to deploy before @jambonz/siprec-client-utils is bumped.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: drop the siprec re-subscribe, keep the missing teardown

Running the cross-feature-server move against a real two-feature-server cluster
showed the rtpengine subscription survives the REFER on its own: with
resubscribe() removed, the fork still carried the post-transfer conversation
(50 packets/s on both streams, and the words spoken after the move transcribed
straight out of the recording). rtpengine keeps non-offer-answer subscriptions
across an answer, so there was never anything to rebuild - the earlier commit
was fixing a fault that does not exist.

What does fail, and what this branch still fixes, is the teardown: the destroy
handler installed on the transferred leg never stopped the recording, so the
recorder was left without a BYE.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-07 18:42:20 +01:00
Dave HortonandClaude Fable 5 b7b707cc2e fix: complete autoscale scale-in reliably and reject new INVITEs while draining (#248)
Scale-in completion never happened: app.js holds the placeholder
Emitter that autoscale-manager returns synchronously (the real
SnsNotifier replaces it later inside an async IIFE), so the completion
poller never saw operationalState change; it also called the
nonexistent scaleIn() rather than completeScaleIn().  Instances in
Terminating:Wait therefore always burned the full lifecycle hook
heartbeat timeout.

In addition, nothing consumed dryUpCalls: a draining SBC kept accepting
new INVITEs sent directly to its public address right up until
termination.

Changes:
- complete the scale-in from within the ScaleIn handler in
  autoscale-manager, where the real notifier is in scope
- while draining, reject new INVITEs with 503 so senders fail over to
  another SBC (INVITE with Replaces is allowed through since it targets
  a call already in progress here)
- a server may run several sbc-inbound and sbc-outbound processes, and
  completing the hook when only this process is idle would terminate
  the instance while sibling processes still have calls; each process
  now reports its call count to redis (lib/call-count-reporter.js, with
  a companion change in sbc-outbound) and the draining process
  completes only when the server-wide count is zero on two consecutive
  checks, falling back to its own count if no reports are present

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
2026-08-31 08:54:19 -04:00
Dave Horton ccbebc015f 0.9.9 2026-08-20 09:16:38 -04:00
Dave Horton 320a30e328 0.9.8 2026-07-20 08:42:25 -04:00
Dave Horton cd2fde360d 0.9.7 2026-07-20 08:39:53 -04:00
Hoan Luu Huu e353b751c0 fixed wrongly identify 3pcc call (#245)
* fixed wrongly identify 3pcc call

* update drachtio 5.0.27
2026-07-11 08:31:59 -04:00
Hoan Luu Huu 53b46f7e51 update rtcp-mux to default for srtp (#246) 2026-07-09 07:09:53 -04:00
Hoan Luu Huu 6608286e8c fixed: don't log carrier register_credential to log (#244)
* fixed: don't log carrier register_credential to log

* upgrade realtimedb helper 0.8.21
2026-06-19 07:08:52 -04:00
Dave HortonandClaude Opus 4.5 6ddfbc9373 fix error handler middleware parameter order (#243)
The error handler was using (req, res, next, err) but drachtio-srf
error middleware expects (err, req, res, next). This caused
"res.send is not a function" errors when the handler was invoked.

Co-authored-by: Claude Opus 4.5 <noreply@anthropic.com>
2026-05-21 15:20:17 -04:00
Dave Horton cd48675499 curtail massive log msg (#242) 2026-04-23 07:35:27 -04:00
Hoan Luu Huu cecdbdccef update drachtio srf 5.0.21 (#241) 2026-04-13 21:24:28 -04:00
Dave Horton 6586919c86 update deps 2026-03-30 21:07:53 -04:00
Dave Horton c919438af5 update to drachtio-srf latest (#239) 2026-03-05 09:24:49 -05:00
Dave Horton ab39525467 update drachtio-srf (#236) 2026-02-09 10:28:42 -05:00
Hoan Luu Huu 61a66ab181 update drachtio srf version (#231) 2026-01-22 08:13:11 -05:00
Dave Horton 48dd7ebfcd Parallelize independent exact IP and CIDR gateway database queries to reduce latency in inbound call routing (#234) 2026-01-18 09:34:48 -05:00
Dave Horton 678fe9d9a8 Fix/sql query optimize (#233)
* Add STRAIGHT_JOIN to CIDR gateway query to prevent MySQL optimizer from choosing inefficient full table scan on voip_carriers table.

* update tests with latest schema

* fix test data

* security issues

* update workflow actions
2026-01-16 08:53:22 -05:00
Dave Horton 39bd65fb97 improve slow sql query and fix support for readonly endpoints (#232) 2026-01-15 08:47:04 -05:00
Sam Machin 3f0fd794ce Recording update (#229)
* new hasRecording flag on call info for setting recording URL in cdr

* lint

* use nullish coalescing for null response
2026-01-02 10:29:00 -05:00
Dave Horton 876cb393f9 fix vulnerability 2025-12-09 09:58:46 -05:00
Sam Machin fa2259d6e5 fix build (#228)
* update dependenices

* Update package-lock.json

* fix buold
2025-12-08 10:43:47 -05:00
Sam Machin 2405190dcc update dependenices (#227) 2025-12-08 09:22:27 -05:00
Hoan Luu Huu e9a5921b19 allow uas leg can send re-invite with outbound gatway credential (#221)
* allow uas leg can send re-invite with outbound gatway credential

* fix cannot get register username/password

* fix cannot get register username/password

* fix cannot get register username/password

* fixed failing test cases
2025-11-24 20:16:28 -06:00
Dave Horton 6efe714d49 fix package lock 2025-11-14 07:36:42 -05:00
Sam Machin 4150454736 remove no SP gateways if Account Gateways (#219) 2025-11-11 08:30:30 -05:00
Anton Voylenko 73eb8e8d17 chore: bump node version (#220) 2025-11-04 18:07:01 -05:00
Hoan Luu Huu 9beba4330a support auth trunk for incoming call (#213)
* support auth trunk for incoming call

* wip

* wip

* wip

* update digest-utils version

* update sql file from api-server

* update sql file from api-server

* wip
2025-10-23 17:00:34 -04:00
14 changed files with 1484 additions and 939 deletions
+4 -4
View File
@@ -15,7 +15,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@v3
uses: actions/checkout@v4
- name: prepare tag
id: prepare_tag
@@ -37,14 +37,14 @@ jobs:
echo "image_id=$IMAGE_ID" >> $GITHUB_OUTPUT
echo "version=$VERSION" >> $GITHUB_OUTPUT
- name: Login to Docker Hub
uses: docker/login-action@v2
- name: Login to Docker Hub
uses: docker/login-action@v3
with:
username: ${{ secrets.DOCKERHUB_USERNAME }}
password: ${{ secrets.DOCKERHUB_TOKEN }}
- name: Build and push Docker image
uses: docker/build-push-action@v4
uses: docker/build-push-action@v6
with:
context: .
push: true
+3 -3
View File
@@ -1,10 +1,10 @@
FROM --platform=linux/amd64 node:20.13.0-alpine3.18 as base
FROM --platform=linux/amd64 node:24-alpine AS base
RUN apk --update --no-cache add --virtual .builds-deps build-base python3
WORKDIR /opt/app/
FROM base as build
FROM base AS build
COPY package.json package-lock.json ./
@@ -18,6 +18,6 @@ COPY --from=build /opt/app /opt/app/
ARG NODE_ENV
ENV NODE_ENV $NODE_ENV
ENV NODE_ENV=$NODE_ENV
CMD [ "node", "app.js" ]
+32 -6
View File
@@ -64,12 +64,21 @@ const {
password: process.env.JAMBONES_MYSQL_PASSWORD,
database: process.env.JAMBONES_MYSQL_DATABASE,
connectionLimit: process.env.JAMBONES_MYSQL_CONNECTION_LIMIT || 10
}, logger);
}, logger, process.env.JAMBONES_MYSQL_WRITE_HOST && process.env.JAMBONES_MYSQL_WRITE_USER &&
process.env.JAMBONES_MYSQL_WRITE_PASSWORD && process.env.JAMBONES_MYSQL_WRITE_DATABASE ? {
host: process.env.JAMBONES_MYSQL_WRITE_HOST,
port: process.env.JAMBONES_MYSQL_WRITE_PORT || 3306,
user: process.env.JAMBONES_MYSQL_WRITE_USER,
password: process.env.JAMBONES_MYSQL_WRITE_PASSWORD,
database: process.env.JAMBONES_MYSQL_WRITE_DATABASE,
connectionLimit: process.env.JAMBONES_MYSQL_CONNECTION_LIMIT || 10
} : null);
const {
client: redisClient,
addKey,
deleteKey,
retrieveKey,
retrieveHash,
createSet,
retrieveSet,
addToSet,
@@ -117,6 +126,7 @@ srf.locals = {...srf.locals,
addKey,
deleteKey,
retrieveKey,
retrieveHash,
createSet,
incrKey,
decrKey,
@@ -130,7 +140,8 @@ const {
wasOriginatedFromCarrier,
getApplicationForDidAndCarrier,
getOutboundGatewayForRefer,
getApplicationBySid
getApplicationBySid,
lookupAuthCarriersForAccountAndSP
} = require('./lib/db-utils')(srf, logger);
srf.locals = {
...srf.locals,
@@ -139,16 +150,29 @@ srf.locals = {
getApplicationForDidAndCarrier,
getOutboundGatewayForRefer,
getFeatureServer: require('./lib/fs-tracking')(srf, logger),
getApplicationBySid
getApplicationBySid,
lookupAuthCarriersForAccountAndSP
};
const activeCallIds = srf.locals.activeCallIds;
/* report our call count to redis so a draining process can count calls across
all sbc-inbound and sbc-outbound processes on this server */
if (!process.env.K8S && 'test' !== process.env.NODE_ENV) {
srf.locals.callCountReporter = require('./lib/call-count-reporter')({
logger,
addKey,
addToSet,
getCount: () => activeCallIds.size
});
}
const {
initLocals,
handleSipRec,
identifyAccount,
checkLimits,
challengeDeviceCalls
challengeDeviceCalls,
identifyAuthTrunk
} = require('./lib/middleware')(srf, logger);
const CallSession = require('./lib/call-session');
@@ -236,7 +260,9 @@ srf.use('invite', [
handleSipRec,
identifyAccount,
checkLimits,
challengeDeviceCalls
challengeDeviceCalls,
// challengeDeviceCalls will detect auth_trunk or device calls, identifyAuthTrunk have to be after that
identifyAuthTrunk
]);
srf.invite((req, res) => {
@@ -255,7 +281,7 @@ srf.invite((req, res) => {
session.connect();
});
srf.use((req, res, next, err) => {
srf.use((err, req, res, next) => {
logger.error(err, 'hit top-level error handler');
res.send(500);
});
+1 -1
View File
@@ -4,7 +4,7 @@
"ICE": "default",
"SDES": "off",
"flags": ["generate mid", "SDES-no", "port latching"],
"rtcp-mux": ["require"]
"rtcp-mux": ["offer"]
},
"teams": {
"transport-protocol": "RTP/SAVP",
+1 -1
View File
@@ -96,7 +96,7 @@ module.exports = [
// Variables
'no-delete-var': 2,
'no-undef': 2,
'no-unused-vars': [2, {args: 'none'}],
'no-unused-vars': [2, {args: 'none', ignoreRestSiblings: true}],
// Node.js and CommonJS
'no-mixed-requires': 2,
+42 -9
View File
@@ -23,23 +23,55 @@ module.exports = (logger) => {
const {srf} = require('..');
const {activeCallIds, removeFromRedis} = srf.locals;
/* reject new INVITEs with 503 so senders fail over to another SBC */
srf.locals.dryUpCalls = true;
/* 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`);
}
/* count calls in progress across all sbc-inbound and sbc-outbound
processes on this server, if they are reporting; otherwise
fall back to counting only our own */
const countServerCalls = async() => {
const reporter = srf.locals.callCountReporter;
if (!reporter) return activeCallIds.size;
const {retrieveSet, retrieveKey} = srf.locals.realtimeDbHelpers;
const keys = await retrieveSet(reporter.setName);
let count = 0;
for (const key of keys) {
count += parseInt(await retrieveKey(key), 10) || 0;
}
return Math.max(count, activeCallIds.size);
};
/* poll until calls have dried up, then complete the scale-in;
require two consecutive zero readings since reported counts
may be up to 15s stale */
let consecutiveZeroCounts = 0;
const timer = setInterval(async() => {
try {
const calls = await countServerCalls();
if (0 === calls) {
if (++consecutiveZeroCounts >= 2) {
clearInterval(timer);
logger.info('scale-in complete now that calls have dried up');
lifecycleEmitter.completeScaleIn();
}
}
else {
consecutiveZeroCounts = 0;
logger.info(`${calls} calls in progress on this server; scale-in will complete when they are done`);
}
} catch (err) {
logger.error({err}, 'Error counting calls in progress during scale-in');
}
}, 20000);
})
.on(LifeCycleEvents.StandbyEnter, () => {
lifecycleEmitter.dryUpCalls = true;
const {srf} = require('..');
const {removeFromRedis} = srf.locals;
srf.locals.dryUpCalls = true;
removeFromRedis();
logger.info('AWS enter pending state notification: begin drying up calls');
@@ -48,6 +80,7 @@ module.exports = (logger) => {
lifecycleEmitter.dryUpCalls = false;
const {srf} = require('..');
const {addToRedis} = srf.locals;
srf.locals.dryUpCalls = false;
addToRedis();
logger.info('AWS exit pending state notification: re-enable calls');
+31
View File
@@ -0,0 +1,31 @@
const os = require('os');
/**
* Periodically report this process's count of calls in progress to redis.
* A server may host several sbc-inbound and sbc-outbound processes; when one
* of them handles an autoscale drain it needs to know when the entire server
* has no calls in progress, not just its own process. Each process writes
* its own count under a per-pid key (with a short expiry, so keys from dead
* processes evaporate) and registers that key in a per-host set that the
* draining process can enumerate.
*/
const REPORT_INTERVAL = 15000;
const KEY_EXPIRY_SECS = 120;
module.exports = ({logger, addKey, addToSet, getCount}) => {
const prefix = process.env.JAMBONES_CLUSTER_ID || 'default';
const setName = `${prefix}:call-count-keys:${os.hostname()}`;
const key = `${prefix}:call-count:${os.hostname()}:${process.pid}`;
const report = () => {
addKey(key, `${getCount()}`, KEY_EXPIRY_SECS)
.catch((err) => logger.error({err}, 'call-count-reporter: error writing call count'));
};
addToSet(setName, key)
.catch((err) => logger.error({err}, `call-count-reporter: error adding ${key} to ${setName}`));
setInterval(report, REPORT_INTERVAL);
report();
return {key, setName};
};
+19 -2
View File
@@ -66,6 +66,7 @@ class CallSession extends Emitter {
this.decrKey = req.srf.locals.realtimeDbHelpers.decrKey;
this.addKey = req.srf.locals.realtimeDbHelpers.addKey;
this.retrieveKey = req.srf.locals.realtimeDbHelpers.retrieveKey;
this.retrieveHash = req.srf.locals.realtimeDbHelpers.retrieveHash;
this._mediaPath = MediaPath.FullMedia;
@@ -108,7 +109,11 @@ class CallSession extends Emitter {
async connect() {
const {sdp} = this.req.locals;
const is3pcc = this.req.body?.length === 0;
// use the parsed SDP from middleware (req.locals.sdp), not req.body:
// for multipart (SIPREC) or IWF-originated bodies req.body may be empty
// or not represent the actual offer, which would wrongly send an
// SDP-bearing call down the no-offer 3pcc path.
const is3pcc = !sdp || sdp.length === 0;
this.logger.info(`inbound ${is3pcc ? '3pcc ' : ''}call accepted for routing`);
const engine = this.getRtpEngine();
if (!engine) {
@@ -332,6 +337,15 @@ class CallSession extends Emitter {
},
});
// passing gateway outbound auth that sbc-inbound can send RE-INVITE with authentication process on UAS side
const {gateway} = this.req.locals;
if (gateway && gateway.register_username && gateway.register_password) {
this.logger.debug('passing outbound gateway auth to CallSession for reinvite processing');
uas.auth = {
username: gateway.register_username,
password: gateway.register_password
};
}
// successfully connected
this.logger.info('call connected successfully to feature server');
debug('call connected successfully to feature server');
@@ -482,7 +496,9 @@ class CallSession extends Emitter {
this.req.locals.carrier :
this.req.locals.originator;
const application = await this.srf.locals.getApplicationBySid(application_sid);
const isRecording = this.req.locals.account.record_all_calls || (application && application.record_all_calls);
const {hasRecording} = await this.retrieveHash(`call:${this.account_sid}:${call_sid}`) ?? {};
const isRecording = this.req.locals.account.record_all_calls ||
(application && application.record_all_calls) || hasRecording;
const day = new Date();
let recording_url = `/Accounts/${this.account_sid}/RecentCalls/${call_sid}/record`;
recording_url += `/${day.getFullYear()}/${(day.getMonth() + 1).toString().padStart(2, '0')}`;
@@ -1093,6 +1109,7 @@ Duration=${payload.duration} `
this.rtpEngineResource.destroy();
this.activeCallIds.delete(this.req.get('Call-ID'));
uac.other.destroy();
this._stopRecording();
this.srf.endSession(this.req);
});
+100 -23
View File
@@ -7,7 +7,8 @@ const sqlSelectSPForAccount = 'SELECT service_provider_sid FROM accounts WHERE a
const sqlSelectAllCarriersForAccountByRealm =
`SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.account_sid,
vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask, sg.pad_crypto
vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask, sg.pad_crypto,
vc.register_username, vc.register_password
FROM sip_gateways sg, voip_carriers vc, accounts acc
WHERE acc.sip_realm = ?
AND vc.account_sid = acc.account_sid
@@ -17,7 +18,8 @@ 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, sg.pad_crypto
vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask, sg.pad_crypto,
vc.register_username, vc.register_password
FROM sip_gateways sg, voip_carriers vc, accounts acc
WHERE acc.sip_realm = ?
AND vc.service_provider_sid = acc.service_provider_sid
@@ -26,14 +28,30 @@ AND vc.is_active = 1
AND sg.inbound = 1
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.account_sid, vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask, sg.pad_crypto
const sqlSelectExactGatewayForSP =
`SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.service_provider_sid,
vc.account_sid, vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask, sg.pad_crypto,
vc.register_username, vc.register_password
FROM sip_gateways sg, voip_carriers vc
WHERE sg.voip_carrier_sid = vc.voip_carrier_sid
AND vc.service_provider_sid IS NOT NULL
AND vc.is_active = 1
AND sg.inbound = 1`;
WHERE sg.voip_carrier_sid = vc.voip_carrier_sid
AND vc.service_provider_sid IS NOT NULL
AND vc.is_active = 1
AND sg.inbound = 1
AND sg.netmask = 32
AND sg.ipv4 = ?
ORDER BY vc.account_sid IS NOT NULL DESC`;
const sqlSelectCIDRGatewaysForSP =
`SELECT STRAIGHT_JOIN sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.service_provider_sid,
vc.account_sid, vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask, sg.pad_crypto,
vc.register_username, vc.register_password
FROM sip_gateways sg, voip_carriers vc
WHERE sg.voip_carrier_sid = vc.voip_carrier_sid
AND vc.service_provider_sid IS NOT NULL
AND vc.is_active = 1
AND sg.inbound = 1
AND sg.netmask < 32
ORDER BY sg.netmask DESC`;
const sqlAccountByRealm = 'SELECT * from accounts WHERE sip_realm = ? AND is_active = 1';
const sqlAccountBySid = 'SELECT * from accounts WHERE account_sid = ?';
@@ -57,7 +75,8 @@ AND outbound = 1`;
const sqlSelectCarrierRequiringRegistration = `
SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.service_provider_sid, vc.account_sid,
vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask, sg.pad_crypto
vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask, sg.pad_crypto,
vc.register_username, vc.register_password
FROM sip_gateways sg, voip_carriers vc
WHERE sg.voip_carrier_sid = vc.voip_carrier_sid
AND vc.requires_register = 1
@@ -65,9 +84,20 @@ AND vc.is_active = 1
AND vc.register_sip_realm = ?
AND vc.register_username = ?`;
const sqlSelectAuthCarriersForAccountAndSP = `
SELECT * FROM voip_carriers
WHERE trunk_type = 'auth'
AND is_active = 1
AND (
(account_sid = ?)
OR
(service_provider_sid = ? AND account_sid IS NULL)
)`;
const sqlSelectGatewaysByVoipCarrierSids = `
SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.service_provider_sid,
vc.account_sid, vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask, sg.pad_crypto
vc.account_sid, vc.application_sid, sg.inbound, sg.outbound, sg.is_active, sg.ipv4, sg.netmask, sg.pad_crypto,
vc.register_username, vc.register_password
FROM sip_gateways sg, voip_carriers vc
WHERE sg.voip_carrier_sid IN (?)
AND sg.voip_carrier_sid = vc.voip_carrier_sid
@@ -284,7 +314,7 @@ module.exports = (srf, logger) => {
/* get all the carriers and gateways for the account owning this sip realm */
const [gwAcc] = await pp.query(sqlSelectAllCarriersForAccountByRealm, [uri.host]);
const [gwSP] = gwAcc.length ? [[]] : await pp.query(sqlSelectAllCarriersForSPByRealm, uri.host);
const [gwSP] = await pp.query(sqlSelectAllCarriersForSPByRealm, uri.host);
const gw = gwAcc
.concat(gwSP)
.sort((a, b) => b.netmask - a.netmask);
@@ -295,7 +325,9 @@ module.exports = (srf, logger) => {
name: gw.name,
service_provider_sid: gw.service_provider_sid,
account_sid: gw.account_sid,
application_sid: gw.application_sid
application_sid: gw.application_sid,
register_username: gw.register_username,
register_password: gw.register_password
};
});
/* remove duplicates, winnow down to voip_carriers, not gateways */
@@ -316,7 +348,8 @@ module.exports = (srf, logger) => {
service_provider_sid: gw.service_provider_sid,
account_sid: gw.account_sid,
application_sid: gw.application_sid,
pad_crypto: gw.pad_crypto
register_username: gw.register_username,
register_password: gw.register_password
};
});
/* remove duplicates */
@@ -361,8 +394,13 @@ module.exports = (srf, logger) => {
}
}
if (r.length > 1) {
logger.info({r},
'multiple carriers with the same gateway have the same number provisioned for the same account'
const all = r.map(({account_sid, voip_carrier_sid}) => ({account_sid, voip_carrier_sid}));
logger.info({
number: r[0].number,
total: all.length,
matches: all.slice(0, 5)
},
'multiple carriers with the same gateway have the same number provisioned for the same account'
+ ' -- cannot determine which one to use');
return {
fromCarrier: true,
@@ -437,10 +475,19 @@ module.exports = (srf, logger) => {
}
/* find all carrier entries that have an inbound gateway matching the source IP */
const [gw] = await pp.query(sqlSelectAllGatewaysForSP);
/* Query both exact IP matches AND CIDR ranges in parallel to handle the case where
multiple accounts have configured the same carrier with different netmasks.
The phone number lookup will disambiguate which account owns the call. */
const [[gwExact], [gwCidr]] = await Promise.all([
pp.query(sqlSelectExactGatewayForSP, [req.source_address]),
pp.query(sqlSelectCIDRGatewaysForSP)
]);
/* Merge both result sets - exact matches first, then CIDR ranges (already sorted by netmask DESC) */
const gw = [...gwExact, ...gwCidr];
//logger.debug({gw}, `checking gateways for source address ${req.source_address}`);
let matches = gw
.sort((a, b) => b.netmask - a.netmask)
.filter(gatewayMatchesSourceAddress.bind(null, logger, req.source_address))
.map((gw) => {
return {
@@ -449,7 +496,9 @@ module.exports = (srf, logger) => {
service_provider_sid: gw.service_provider_sid,
account_sid: gw.account_sid,
application_sid: gw.application_sid,
pad_crypto: gw.pad_crypto
pad_crypto: gw.pad_crypto,
register_username: gw.register_username,
register_password: gw.register_password
};
});
/* remove duplicates, winnow down to voip_carriers, not gateways */
@@ -469,7 +518,9 @@ module.exports = (srf, logger) => {
service_provider_sid: gw.service_provider_sid,
account_sid: gw.account_sid,
application_sid: gw.application_sid,
pad_crypto: gw.pad_crypto
pad_crypto: gw.pad_crypto,
register_username: gw.register_username,
register_password: gw.register_password
};
});
/* remove duplicates */
@@ -555,8 +606,13 @@ module.exports = (srf, logger) => {
}
}
else if (r.length > 1) {
logger.info({r},
'multiple accounts have added this carrier with default routing -- cannot determine which to use');
const all = r.map(({account_sid, voip_carrier_sid}) => ({account_sid, voip_carrier_sid}));
logger.info({
number: r[0].number,
total: all.length,
matches: all.slice(0, 5)
},
'multiple accounts have added this carrier with default routing -- cannot determine which to use');
return {
fromCarrier: true,
error: 'Multiple accounts are attempting to route the same phone number from the same carrier'
@@ -579,12 +635,33 @@ module.exports = (srf, logger) => {
return failure;
};
/**
* Retrieves voip_carriers with trunk_type 'auth' that belong to either:
* 1. The specified account (account_sid matches), OR
* 2. The service provider but with null account_sid (shared across service provider)
*
* @param {string} account_sid - The SID of the account
* @param {string} service_provider_sid - The SID of the service provider
* @returns {Promise<Array>} Array of voip_carrier records matching the criteria
* @throws {Error} Database errors or other unexpected errors
*/
const lookupAuthCarriersForAccountAndSP = async(account_sid, service_provider_sid) => {
try {
const [rows] = await pp.query(sqlSelectAuthCarriersForAccountAndSP, [account_sid, service_provider_sid]);
return rows;
} catch (err) {
logger.error({err, account_sid, service_provider_sid}, 'lookupAuthCarriersForAccountAndSP');
throw err;
}
};
return {
wasOriginatedFromCarrier,
getApplicationForDidAndCarrier,
getApplicationForDidAndCarriers,
getOutboundGatewayForRefer,
getSPForAccount,
getApplicationBySid
getApplicationBySid,
lookupAuthCarriersForAccountAndSP
};
};
+51 -3
View File
@@ -28,12 +28,21 @@ module.exports = function(srf, logger) {
lookupAccountBySipRealm,
lookupAccountBySid,
lookupAccountCapacitiesBySid,
queryCallLimits
queryCallLimits,
} = srf.locals.dbHelpers;
const {stats, writeCdrs} = srf.locals;
const {stats, writeCdrs, lookupAuthCarriersForAccountAndSP, getApplicationForDidAndCarrier} = srf.locals;
const initLocals = (req, res, next) => {
const callId = req.get('Call-ID');
/* if we are drying up calls prior to scale-in, reject new INVITEs so the
sender fails over to another SBC; allow INVITE with Replaces through
since it targets a call already in progress on this server */
if (srf.locals.dryUpCalls && !req.has('Replaces')) {
logger.info({callId}, 'rejecting INVITE with 503 as we are drying up calls before scale-in');
return res.send(503);
}
req.locals = req.locals || {callId};
req.locals.nudge = 0;
@@ -130,7 +139,8 @@ module.exports = function(srf, logger) {
logger.info('identifyAccount: rejecting call from carrier because DID has not been provisioned');
return res.send(404, 'Number Not Provisioned');
}
logger.info({gateway}, 'identifyAccount: incoming call from gateway');
const {register_password, ...gatewayForLog} = gateway;
logger.info({gateway: gatewayForLog}, 'identifyAccount: incoming call from gateway');
const appSidHeader = req.get('x-application-sid');
if (appSidHeader && appSidHeader == application_sid) {
logger.info({callId}, 'Loop Detected, x-application-sid header on incoming call matches applicationSid');
@@ -191,6 +201,10 @@ module.exports = function(srf, logger) {
res.send(404);
return req.srf.endSession(req);
}
const auth_trunks = await lookupAuthCarriersForAccountAndSP(
account.account_sid,
account.service_provider_sid
);
/* if this is a dedicated SBC (static IP) only take calls for that account's sip realm */
if (process.env.SBC_ACCOUNT_SID && account.account_sid !== process.env.SBC_ACCOUNT_SID) {
@@ -214,6 +228,7 @@ module.exports = function(srf, logger) {
registration_hook_username: account.registration_hook.username,
registration_hook_password: account.registration_hook.password
}),
...(auth_trunks?.length && {auth_trunks}),
...req.locals
};
}
@@ -365,6 +380,38 @@ module.exports = function(srf, logger) {
}
};
const identifyAuthTrunk = async(req, res, next) => {
try {
if (req.authorization) {
const {grant} = req.authorization;
if (grant && grant.status === 'ok' && grant.auth_trunk) {
// we have successfully authenticated the call for an auth_trunk
const application_sid = await getApplicationForDidAndCarrier(req, grant.auth_trunk.voip_carrier_sid);
req.locals = {
...req.locals,
originator: 'trunk',
carrier: grant.auth_trunk.name,
gateway: grant.auth_trunk,
voip_carrier_sid: grant.auth_trunk.voip_carrier_sid,
application_sid: application_sid || grant.auth_trunk.application_sid,
};
// as call from auth carrier, clean req.authorization that impact on legacy logic for authenticated user
delete req.authorization;
logger.debug({callId: req.locals.callId, auth_trunk: grant.auth_trunk.name},
'identifyAuthTrunk: call authenticated for auth trunk');
}
}
next();
} catch (err) {
stats.increment('sbc.terminations', ['sipStatus:500']);
logger.error(err, `${req.get('Call-ID')} Error challenging auth trunk`);
res.send(500);
req.srf.endSession(req);
}
};
const challengeDeviceCalls = async(req, res, next) => {
try {
/* TODO: check if this is a gateway that we have an ACL for */
@@ -383,6 +430,7 @@ module.exports = function(srf, logger) {
handleSipRec,
challengeDeviceCalls,
identifyAccount,
identifyAuthTrunk,
checkLimits
};
};
+1165 -856
View File
File diff suppressed because it is too large Load Diff
+7 -7
View File
@@ -1,9 +1,9 @@
{
"name": "sbc-inbound",
"version": "0.9.5",
"version": "0.9.11",
"main": "app.js",
"engines": {
"node": ">= 18.0.0"
"node": ">= 20.0.0"
},
"keywords": [
"sip",
@@ -30,19 +30,19 @@
"@aws-sdk/client-sns": "^3.549.0",
"@babel/helpers": "^7.26.10",
"@jambonz/db-helpers": "^0.9.18",
"@jambonz/digest-utils": "^0.0.8",
"@jambonz/digest-utils": "^0.0.9",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/realtimedb-helpers": "^0.8.18",
"@jambonz/realtimedb-helpers": "^0.8.21",
"@jambonz/rtpengine-utils": "^0.4.4",
"@jambonz/siprec-client-utils": "^0.2.10",
"@jambonz/stats-collector": "^0.1.10",
"@jambonz/time-series": "^0.2.10",
"bent": "^7.3.12",
"cidr-matcher": "^2.1.1",
"debug": "^4.3.4",
"debug": "^4.4.3",
"drachtio-fn-b2b-sugar": "0.2.1",
"drachtio-srf": "^5.0.5",
"express": "^4.19.2",
"drachtio-srf": "^5.0.27",
"express": "^4.21.2",
"pino": "^10.1.0",
"verify-aws-sns-signature": "^0.1.0",
"xml2js": "^0.6.2"
+27 -23
View File
@@ -14,6 +14,8 @@ DROP TABLE IF EXISTS beta_invite_codes;
DROP TABLE IF EXISTS call_routes;
DROP TABLE IF EXISTS clients;
DROP TABLE IF EXISTS dns_records;
DROP TABLE IF EXISTS lcr;
@@ -66,8 +68,6 @@ DROP TABLE IF EXISTS phone_numbers;
DROP TABLE IF EXISTS sip_gateways;
DROP TABLE IF EXISTS clients;
DROP TABLE IF EXISTS voip_carriers;
DROP TABLE IF EXISTS accounts;
@@ -132,6 +132,19 @@ application_sid CHAR(36) NOT NULL,
PRIMARY KEY (call_route_sid)
) COMMENT='a regex-based pattern match for call routing';
CREATE TABLE clients
(
client_sid CHAR(36) NOT NULL UNIQUE ,
account_sid CHAR(36) NOT NULL,
is_active BOOLEAN NOT NULL DEFAULT 1,
username VARCHAR(64),
password VARCHAR(1024),
allow_direct_app_calling BOOLEAN NOT NULL DEFAULT 1,
allow_direct_queue_calling BOOLEAN NOT NULL DEFAULT 1,
allow_direct_user_calling BOOLEAN NOT NULL DEFAULT 1,
PRIMARY KEY (client_sid)
);
CREATE TABLE dns_records
(
dns_record_sid CHAR(36) NOT NULL UNIQUE ,
@@ -191,6 +204,7 @@ tech_prefix VARCHAR(16) COMMENT 'tech prefix to prepend to outbound calls to thi
inbound_auth_username VARCHAR(64),
inbound_auth_password VARCHAR(64),
diversion VARCHAR(32),
trunk_type ENUM('static_ip','auth','reg') NOT NULL DEFAULT 'static_ip',
PRIMARY KEY (predefined_carrier_sid)
);
@@ -405,7 +419,7 @@ register_public_ip_in_contact BOOLEAN NOT NULL DEFAULT false,
register_status VARCHAR(4096),
dtmf_type ENUM('rfc2833','tones','info') NOT NULL DEFAULT 'rfc2833',
outbound_sip_proxy VARCHAR(255),
trunk_type ENUM('static-ip','auth','registration') NOT NULL DEFAULT 'static-ip',
trunk_type ENUM('static_ip','auth','reg') NOT NULL DEFAULT 'static_ip',
PRIMARY KEY (voip_carrier_sid)
) COMMENT='A Carrier or customer PBX that can send or receive calls';
@@ -479,20 +493,6 @@ password VARCHAR(255),
PRIMARY KEY (webhook_sid)
) COMMENT='An HTTP callback';
CREATE TABLE clients
(
client_sid CHAR(36) NOT NULL UNIQUE ,
account_sid CHAR(36) NOT NULL,
is_active BOOLEAN NOT NULL DEFAULT 1,
username VARCHAR(64),
password VARCHAR(1024),
allow_direct_app_calling BOOLEAN NOT NULL DEFAULT 1,
allow_direct_queue_calling BOOLEAN NOT NULL DEFAULT 1,
allow_direct_user_calling BOOLEAN NOT NULL DEFAULT 1,
voip_carrier_sid CHAR(36),
PRIMARY KEY (client_sid)
);
CREATE TABLE applications
(
application_sid CHAR(36) NOT NULL UNIQUE ,
@@ -505,7 +505,7 @@ messaging_hook_sid CHAR(36) COMMENT 'webhook to call for inbound SMS/MMS ',
app_json TEXT,
speech_synthesis_vendor VARCHAR(64) NOT NULL DEFAULT 'google',
speech_synthesis_language VARCHAR(12) NOT NULL DEFAULT 'en-US',
speech_synthesis_voice VARCHAR(256),
speech_synthesis_voice VARCHAR(256) DEFAULT 'en-US-Standard-C',
speech_synthesis_label VARCHAR(64),
speech_recognizer_vendor VARCHAR(64) NOT NULL DEFAULT 'google',
speech_recognizer_language VARCHAR(64) NOT NULL DEFAULT 'en-US',
@@ -582,6 +582,9 @@ ALTER TABLE call_routes ADD FOREIGN KEY account_sid_idxfk_3 (account_sid) REFERE
ALTER TABLE call_routes ADD FOREIGN KEY application_sid_idxfk (application_sid) REFERENCES applications (application_sid);
CREATE INDEX client_sid_idx ON clients (client_sid);
ALTER TABLE clients ADD CONSTRAINT account_sid_idxfk_13 FOREIGN KEY account_sid_idxfk_13 (account_sid) REFERENCES accounts (account_sid);
CREATE INDEX dns_record_sid_idx ON dns_records (dns_record_sid);
ALTER TABLE dns_records ADD FOREIGN KEY account_sid_idxfk_4 (account_sid) REFERENCES accounts (account_sid);
@@ -702,6 +705,12 @@ ALTER TABLE phone_numbers ADD FOREIGN KEY service_provider_sid_idxfk_8 (service_
CREATE INDEX sip_gateway_idx_hostport ON sip_gateways (ipv4,port);
CREATE INDEX idx_sip_gateways_inbound_carrier ON sip_gateways (inbound,voip_carrier_sid);
CREATE INDEX idx_sip_gateways_inbound_lookup ON sip_gateways (inbound,netmask,ipv4);
CREATE INDEX idx_sip_gateways_inbound_netmask ON sip_gateways (inbound,netmask);
CREATE INDEX voip_carrier_sid_idx ON sip_gateways (voip_carrier_sid);
ALTER TABLE sip_gateways ADD FOREIGN KEY voip_carrier_sid_idxfk_2 (voip_carrier_sid) REFERENCES voip_carriers (voip_carrier_sid);
@@ -710,11 +719,6 @@ ALTER TABLE lcr_carrier_set_entry ADD FOREIGN KEY lcr_route_sid_idxfk (lcr_route
ALTER TABLE lcr_carrier_set_entry ADD FOREIGN KEY voip_carrier_sid_idxfk_3 (voip_carrier_sid) REFERENCES voip_carriers (voip_carrier_sid);
CREATE INDEX webhook_sid_idx ON webhooks (webhook_sid);
CREATE INDEX client_sid_idx ON clients (client_sid);
ALTER TABLE clients ADD CONSTRAINT account_sid_idxfk_13 FOREIGN KEY account_sid_idxfk_13 (account_sid) REFERENCES accounts (account_sid);
ALTER TABLE clients ADD FOREIGN KEY voip_carrier_sid_idxfk_4 (voip_carrier_sid) REFERENCES voip_carriers (voip_carrier_sid);
CREATE UNIQUE INDEX applications_idx_name ON applications (account_sid,name);
CREATE INDEX application_sid_idx ON applications (application_sid);
+1 -1
View File
@@ -129,7 +129,7 @@ values ('acct-100', 'Account 100', '3f35518f-5a0d-4c2e-90a5-2407bb3b36f0', 'foob
insert into voip_carriers (voip_carrier_sid, name, account_sid, service_provider_sid, trunk_type,
requires_register, register_username, register_sip_realm, register_password, is_active)
values ('4a7d1c8e-5f2b-4d9a-8e3c-6b5a9f1e4c7d', 'test-registration-trunk', 'ed649e33-e771-403a-8c99-1780eabbc803',
'3f35518f-5a0d-4c2e-90a5-2407bb3b36f0', 'registration', true, 'testuser',
'3f35518f-5a0d-4c2e-90a5-2407bb3b36f0', 'reg', true, 'testuser',
'sip.carrier.example.com', 'testpass', true);
-- sip_gateway for outbound only (inbound will use ephemeral gateway from Redis)