Compare commits

..
Author SHA1 Message Date
rhonda hollis a5ea0740a9 refactor
- remove Redis
- handle proxing of non-success response
2026-05-15 17:28:17 -07:00
rhonda hollis 79b2a99fe8 ensure application_sid & traceId are on CDR 2026-05-14 09:14:43 -07:00
9 changed files with 86 additions and 128 deletions
+1 -12
View File
@@ -155,17 +155,6 @@ srf.locals = {
};
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,
@@ -281,7 +270,7 @@ srf.invite((req, res) => {
session.connect();
});
srf.use((err, req, res, next) => {
srf.use((req, res, next, err) => {
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": ["offer"]
"rtcp-mux": ["require"]
},
"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', ignoreRestSiblings: true}],
'no-unused-vars': [2, {args: 'none'}],
// Node.js and CommonJS
'no-mixed-requires': 2,
+9 -42
View File
@@ -23,55 +23,23 @@ 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();
/* 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);
/* 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;
srf.locals.dryUpCalls = true;
removeFromRedis();
logger.info('AWS enter pending state notification: begin drying up calls');
@@ -80,7 +48,6 @@ 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
@@ -1,31 +0,0 @@
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};
};
+34 -6
View File
@@ -19,6 +19,9 @@ const debug = require('debug')('jambonz:sbc-inbound');
const MS_TEAMS_USER_AGENT = 'Microsoft.PSTNHub.SIPProxy';
const MS_TEAMS_SIP_ENDPOINT = 'sip.pstnhub.microsoft.com';
const IMMUTABLE_HEADERS = ['via', 'from', 'to', 'call-id', 'cseq', 'max-forwards', 'content-length'];
const NONCOPYABLE_RESPONSE_HEADERS = [
'via', 'from', 'to', 'call-id', 'cseq', 'contact', 'content-length', 'content-type'
];
/**
* this is to make sure the outgoing From has the number in the incoming From
@@ -109,11 +112,7 @@ class CallSession extends Emitter {
async connect() {
const {sdp} = this.req.locals;
// 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;
const is3pcc = this.req.body?.length === 0;
this.logger.info(`inbound ${is3pcc ? '3pcc ' : ''}call accepted for routing`);
const engine = this.getRtpEngine();
if (!engine) {
@@ -285,6 +284,7 @@ class CallSession extends Emitter {
proxy,
headers,
responseHeaders,
passFailure: false,
proxyRequestHeaders: [
'all',
'-Authorization',
@@ -358,6 +358,19 @@ class CallSession extends Emitter {
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}`);
/* capture trace_id and application_sid from the feature server's error response for the CDR. */
if (this.req.locals.cdr && err.res) {
const trace_id = err.res.get('X-Trace-ID');
if (trace_id) {
this.req.locals.cdr.trace_id = trace_id;
}
const application_sid = err.res.get('X-Application-Sid');
if (application_sid) {
this.req.locals.cdr.application_sid = application_sid;
}
}
this.emit('failed');
}
else if (err.message !== 'call canceled') {
@@ -372,6 +385,22 @@ class CallSession extends Emitter {
.catch((err) => this.logger.error(err, 'Error decrementing call counts'));
}
/* manually proxy the failure response to UAS so that the trace_id/application_sid capture above runs before
res.end fires (which is what triggers the failure-CDR writer). */
if (err.message !== 'call canceled' && !this.res.finalResponseSent) {
if (err instanceof SipError && err.res) {
const headers = {};
Object.keys(err.res.headers || {}).forEach((h) => {
if (!NONCOPYABLE_RESPONSE_HEADERS.includes(h)) headers[h] = err.res.headers[h];
});
this.res.send(err.status, err.reason, {headers});
}
else {
this.res.send(err.status || 500, err.reason);
}
}
this.srf.endSession(this.req);
}
}
@@ -1109,7 +1138,6 @@ Duration=${payload.duration} `
this.rtpEngineResource.destroy();
this.activeCallIds.delete(this.req.get('Call-ID'));
uac.other.destroy();
this._stopRecording();
this.srf.endSession(this.req);
});
+7 -15
View File
@@ -34,15 +34,6 @@ module.exports = function(srf, logger) {
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;
@@ -68,13 +59,15 @@ module.exports = function(srf, logger) {
/* write cdr for non-success response here */
res.once('end', ({status}) => {
if (req.locals.cdr && req.locals.cdr.account_sid && status > 200 && 401 !== status) {
const trunk = ['trunk', 'teams'].includes(req.locals.originator) ? req.locals.carrier : req.locals.originator;
if (req.locals.cdr && req.locals.cdr.account_sid && status > 200 && 401 !== status) {
const trunk = ['trunk', 'teams'].includes(req.locals.originator) ?
req.locals.carrier : req.locals.originator;
writeCdrs({...req.locals.cdr,
terminated_at: Date.now(),
termination_reason: status === 487 === status ? 'caller abandoned' : 'failed',
termination_reason: status === 487 ? 'caller abandoned' : 'failed',
sip_status: status,
trunk
trunk,
...(req.locals.application_sid && {application_sid: req.locals.application_sid})
}).catch((err) => logger.error({err}, 'Error writing cdr for call failure'));
}
});
@@ -139,8 +132,7 @@ module.exports = function(srf, logger) {
logger.info('identifyAccount: rejecting call from carrier because DID has not been provisioned');
return res.send(404, 'Number Not Provisioned');
}
const {register_password, ...gatewayForLog} = gateway;
logger.info({gateway: gatewayForLog}, 'identifyAccount: incoming call from gateway');
logger.info({gateway}, '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');
+30 -17
View File
@@ -1,12 +1,12 @@
{
"name": "sbc-inbound",
"version": "0.9.11",
"version": "0.9.6",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "sbc-inbound",
"version": "0.9.11",
"version": "0.9.6",
"license": "MIT",
"dependencies": {
"@aws-sdk/client-auto-scaling": "^3.549.0",
@@ -15,7 +15,7 @@
"@jambonz/db-helpers": "^0.9.18",
"@jambonz/digest-utils": "^0.0.9",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/realtimedb-helpers": "^0.8.21",
"@jambonz/realtimedb-helpers": "^0.8.18",
"@jambonz/rtpengine-utils": "^0.4.4",
"@jambonz/siprec-client-utils": "^0.2.10",
"@jambonz/stats-collector": "^0.1.10",
@@ -24,7 +24,7 @@
"cidr-matcher": "^2.1.1",
"debug": "^4.4.3",
"drachtio-fn-b2b-sugar": "0.2.1",
"drachtio-srf": "^5.0.27",
"drachtio-srf": "^5.0.21",
"express": "^4.21.2",
"pino": "^10.1.0",
"verify-aws-sns-signature": "^0.1.0",
@@ -1538,10 +1538,9 @@
}
},
"node_modules/@jambonz/realtimedb-helpers": {
"version": "0.8.21",
"resolved": "https://registry.npmjs.org/@jambonz/realtimedb-helpers/-/realtimedb-helpers-0.8.21.tgz",
"integrity": "sha512-Zn/Zw14U0L1pB2RlN5i160FwWqphonhQ7p880OzsKOuvg+wmj6qIkxdLjDXo6TSxFPXDHD25s33zfGOT6ZIfsw==",
"license": "MIT",
"version": "0.8.18",
"resolved": "https://registry.npmjs.org/@jambonz/realtimedb-helpers/-/realtimedb-helpers-0.8.18.tgz",
"integrity": "sha512-PfQRsOy/uKSA0ymRAEmTQnLxc1BZVRQWc+5mClLX6oB6vgvM/6YGgiuv5GSu26OtuKMTtS+RAxlFeT05iH/2rg==",
"dependencies": {
"debug": "^4.3.4",
"ioredis": "^5.3.2"
@@ -3089,9 +3088,9 @@
"integrity": "sha512-bMtje8GWVTze+UG6WSGnlUfBaYCuFiApNXl/XxWc+X9uATZiZkV2jIqhf+Y4SY3nS8dZclJIVKd9qwoCa+i+Vw=="
},
"node_modules/drachtio-srf": {
"version": "5.0.27",
"resolved": "https://registry.npmjs.org/drachtio-srf/-/drachtio-srf-5.0.27.tgz",
"integrity": "sha512-UCNITDPJreCurG5reM2TSmWSbT52azPaNqoc42L4nljy0aVuDt0Am9ExpAy5RfsP8ioxSY2VdSJehr1ARRJL+Q==",
"version": "5.0.21",
"resolved": "https://registry.npmjs.org/drachtio-srf/-/drachtio-srf-5.0.21.tgz",
"integrity": "sha512-9hkQ7LURI1ceTMs49Nh2aosAd2v0565eD8Ccho5sSzXEHuq0DE3epjStlThq6LvIGVVAOXuSQGXFmlMWK92R3w==",
"license": "MIT",
"dependencies": {
"debug": "^4.4.3",
@@ -3099,7 +3098,7 @@
"node-noop": "^1.0.0",
"only": "^0.0.2",
"sdp-transform": "^2.15.0",
"short-uuid": "^6.0.3",
"short-uuid": "^5.2.0",
"sip-methods": "^0.3.0",
"sip-status": "^0.1.0",
"utils-merge": "^1.0.1",
@@ -6111,15 +6110,29 @@
}
},
"node_modules/short-uuid": {
"version": "6.0.3",
"resolved": "https://registry.npmjs.org/short-uuid/-/short-uuid-6.0.3.tgz",
"integrity": "sha512-UMZ3rYOoid307EqWsPTnoBUSOB51Fi28vUMPigUsSCmbaSvslyf/SlXZWOn4btj8tokhBTBkPFKVKKlIolEG1w==",
"version": "5.2.0",
"resolved": "https://registry.npmjs.org/short-uuid/-/short-uuid-5.2.0.tgz",
"integrity": "sha512-296/Nzi4DmANh93iYBwT4NoYRJuHnKEzefrkSagQbTH/A6NTaB68hSPDjm5IlbI5dx9FXdmtqPcj6N5H+CPm6w==",
"license": "MIT",
"dependencies": {
"any-base": "^1.1.0"
"any-base": "^1.1.0",
"uuid": "^9.0.1"
},
"engines": {
"node": ">=14.17.0"
"node": ">=14"
}
},
"node_modules/short-uuid/node_modules/uuid": {
"version": "9.0.1",
"resolved": "https://registry.npmjs.org/uuid/-/uuid-9.0.1.tgz",
"integrity": "sha512-b+1eJOlsR9K8HJpow9Ok3fiWOWSIcIzXodvv0rQjVoOVNpWMpxf1wZNpt4y9h10odCNrqnYp1OBzRktckBe3sA==",
"funding": [
"https://github.com/sponsors/broofa",
"https://github.com/sponsors/ctavan"
],
"license": "MIT",
"bin": {
"uuid": "dist/bin/uuid"
}
},
"node_modules/side-channel": {
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "sbc-inbound",
"version": "0.9.11",
"version": "0.9.6",
"main": "app.js",
"engines": {
"node": ">= 20.0.0"
@@ -32,7 +32,7 @@
"@jambonz/db-helpers": "^0.9.18",
"@jambonz/digest-utils": "^0.0.9",
"@jambonz/http-health-check": "^0.0.1",
"@jambonz/realtimedb-helpers": "^0.8.21",
"@jambonz/realtimedb-helpers": "^0.8.18",
"@jambonz/rtpengine-utils": "^0.4.4",
"@jambonz/siprec-client-utils": "^0.2.10",
"@jambonz/stats-collector": "^0.1.10",
@@ -41,7 +41,7 @@
"cidr-matcher": "^2.1.1",
"debug": "^4.4.3",
"drachtio-fn-b2b-sugar": "0.2.1",
"drachtio-srf": "^5.0.27",
"drachtio-srf": "^5.0.21",
"express": "^4.21.2",
"pino": "^10.1.0",
"verify-aws-sns-signature": "^0.1.0",