Compare commits

..
8 Commits
13 changed files with 878 additions and 1316 deletions
+4 -4
View File
@@ -15,7 +15,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@v4
uses: actions/checkout@v3
- 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@v3
- name: Login to Docker Hub
uses: docker/login-action@v2
with:
username: ${{ secrets.DOCKERHUB_USERNAME }}
password: ${{ secrets.DOCKERHUB_TOKEN }}
- name: Build and push Docker image
uses: docker/build-push-action@v6
uses: docker/build-push-action@v4
with:
context: .
push: true
+2 -23
View File
@@ -64,21 +64,12 @@ const {
password: process.env.JAMBONES_MYSQL_PASSWORD,
database: process.env.JAMBONES_MYSQL_DATABASE,
connectionLimit: process.env.JAMBONES_MYSQL_CONNECTION_LIMIT || 10
}, 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);
}, logger);
const {
client: redisClient,
addKey,
deleteKey,
retrieveKey,
retrieveHash,
createSet,
retrieveSet,
addToSet,
@@ -126,7 +117,6 @@ srf.locals = {...srf.locals,
addKey,
deleteKey,
retrieveKey,
retrieveHash,
createSet,
incrKey,
decrKey,
@@ -155,17 +145,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 +260,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};
};
+2 -10
View File
@@ -66,7 +66,6 @@ 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;
@@ -109,11 +108,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) {
@@ -496,9 +491,7 @@ class CallSession extends Emitter {
this.req.locals.carrier :
this.req.locals.originator;
const application = await this.srf.locals.getApplicationBySid(application_sid);
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 isRecording = this.req.locals.account.record_all_calls || (application && application.record_all_calls);
const day = new Date();
let recording_url = `/Accounts/${this.account_sid}/RecentCalls/${call_sid}/record`;
recording_url += `/${day.getFullYear()}/${(day.getMonth() + 1).toString().padStart(2, '0')}`;
@@ -1109,7 +1102,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);
});
+12 -46
View File
@@ -28,30 +28,15 @@ AND vc.is_active = 1
AND sg.inbound = 1
AND sg.voip_carrier_sid = vc.voip_carrier_sid`;
const sqlSelectExactGatewayForSP =
`SELECT sg.sip_gateway_sid, sg.voip_carrier_sid, vc.name, vc.service_provider_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,
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
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`;
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`;
const sqlAccountByRealm = 'SELECT * from accounts WHERE sip_realm = ? AND is_active = 1';
const sqlAccountBySid = 'SELECT * from accounts WHERE account_sid = ?';
@@ -394,13 +379,8 @@ module.exports = (srf, logger) => {
}
}
if (r.length > 1) {
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'
logger.info({r},
'multiple carriers with the same gateway have the same number provisioned for the same account'
+ ' -- cannot determine which one to use');
return {
fromCarrier: true,
@@ -475,19 +455,10 @@ module.exports = (srf, logger) => {
}
/* find all carrier entries that have an inbound gateway matching the source IP */
/* 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];
const [gw] = await pp.query(sqlSelectAllGatewaysForSP);
//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 {
@@ -606,13 +577,8 @@ module.exports = (srf, logger) => {
}
}
else if (r.length > 1) {
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');
logger.info({r},
'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'
+1 -11
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;
@@ -139,8 +130,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');
+817 -1114
View File
File diff suppressed because it is too large Load Diff
+4 -4
View File
@@ -1,6 +1,6 @@
{
"name": "sbc-inbound",
"version": "0.9.11",
"version": "0.9.5",
"main": "app.js",
"engines": {
"node": ">= 20.0.0"
@@ -30,9 +30,9 @@
"@aws-sdk/client-sns": "^3.549.0",
"@babel/helpers": "^7.26.10",
"@jambonz/db-helpers": "^0.9.18",
"@jambonz/digest-utils": "^0.0.9",
"@jambonz/digest-utils": "^0.0.8",
"@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.5",
"express": "^4.21.2",
"pino": "^10.1.0",
"verify-aws-sns-signature": "^0.1.0",
+24 -28
View File
@@ -14,8 +14,6 @@ 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;
@@ -68,6 +66,8 @@ 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,19 +132,6 @@ 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 ,
@@ -204,7 +191,6 @@ 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)
);
@@ -419,7 +405,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','reg') NOT NULL DEFAULT 'static_ip',
trunk_type ENUM('static-ip','auth','registration') NOT NULL DEFAULT 'static-ip',
PRIMARY KEY (voip_carrier_sid)
) COMMENT='A Carrier or customer PBX that can send or receive calls';
@@ -493,6 +479,20 @@ 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) DEFAULT 'en-US-Standard-C',
speech_synthesis_voice VARCHAR(256),
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,9 +582,6 @@ 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);
@@ -705,12 +702,6 @@ 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);
@@ -719,6 +710,11 @@ 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);
@@ -752,4 +748,4 @@ ALTER TABLE accounts ADD FOREIGN KEY device_calling_application_sid_idxfk (devic
ALTER TABLE accounts ADD FOREIGN KEY siprec_hook_sid_idxfk (siprec_hook_sid) REFERENCES applications (application_sid);
SET FOREIGN_KEY_CHECKS=0;
SET FOREIGN_KEY_CHECKS=0;
+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', 'reg', true, 'testuser',
'3f35518f-5a0d-4c2e-90a5-2407bb3b36f0', 'registration', true, 'testuser',
'sip.carrier.example.com', 'testpass', true);
-- sip_gateway for outbound only (inbound will use ephemeral gateway from Redis)