mirror of
https://github.com/jambonz/sbc-sip-sidecar.git
synced 2026-10-04 02:04:20 +00:00
The Call-ID of outbound registrations was the sip_gateway_sid, the same on every SBC. When the regbot role moved to the other SBC, the registrar saw a refresh of an existing binding from a different source address and Contact. Some registrars 200 such a refresh without updating their routing, so inbound calls to the registered trunk fail with 404 until the binding is recreated. The Call-ID is now sip_gateway_sid@<sending SBC public IP>: stable across refreshes and restarts of one SBC, new when the role moves, so a move looks like a new registration. register_status also records the sending SBC as sbcAddress. With AWS_LIFECYCLE_DRAIN enabled the sidecar polls IMDS autoscaling/target-lifecycle-state (the signal inbound drains on; detection only, inbound completes the lifecycle hook). When the instance is being scaled in, the regbot holder releases the lease while still running instead of after the instance is gone; until now the draining SBC kept the registrations, so carriers kept sending registration-trunk calls to an SBC that answers new INVITEs with 503. It never claims the role back. Once another SBC has claimed it, the draining SBC un-REGISTERs (Expires: 0) the bindings whose Contact carries its own IP. Bindings with an AoR or realm Contact are left alone: the new SBC sends the same Contact, and its REGISTER has already replaced ours. Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
511 lines
19 KiB
JavaScript
511 lines
19 KiB
JavaScript
const debug = require('debug')('jambonz:sbc-registrar');
|
|
const crypto = require('crypto');
|
|
const {
|
|
JAMBONES_CLUSTER_ID,
|
|
JAMBONES_REGBOT_BATCH_SLEEP_MS,
|
|
JAMBONES_REGBOT_BATCH_SIZE,
|
|
AWS_LIFECYCLE_DRAIN,
|
|
} = require('./config');
|
|
const short = require('short-uuid');
|
|
const Regbot = require('./regbot');
|
|
const { sleepFor } = require('./utils');
|
|
|
|
const MAX_INITIAL_DELAY = 15;
|
|
const REGBOT_STATUS_CHECK_INTERVAL = 60;
|
|
/* scale-in handoff: how often to look for a successor, how long to wait for one (a successor
|
|
claims within one status check interval), and how far to trail its registrations */
|
|
const HANDOFF_POLL_MS = 5000;
|
|
const HANDOFF_TIMEOUT_MS = 2 * REGBOT_STATUS_CHECK_INTERVAL * 1000;
|
|
const HANDOFF_GRACE_MS = 10000;
|
|
const regbotKey = `${(JAMBONES_CLUSTER_ID || 'default')}:regbot-token`;
|
|
const waitFor = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
|
|
let initialized = false;
|
|
let rebuildInProgress = false;
|
|
|
|
const regbots = [];
|
|
let carriersHash = '';
|
|
let gatewaysHash = '';
|
|
|
|
function computeHash(arr) {
|
|
const h = crypto.createHash('md5');
|
|
for (const item of arr) h.update(JSON.stringify(item));
|
|
return h.digest('hex');
|
|
}
|
|
|
|
const getCountSuccessfulRegbots = () => regbots.filter((rb) => rb.status === 'registered').length;
|
|
|
|
function pickRelevantCarrierProperties(c) {
|
|
return {
|
|
voip_carrier_sid: c.voip_carrier_sid,
|
|
requires_register: c.requires_register,
|
|
is_active: c.is_active,
|
|
register_username: c.register_username,
|
|
register_password: c.register_password,
|
|
register_sip_realm: c.register_sip_realm,
|
|
register_from_user: c.register_from_user,
|
|
register_from_domain: c.register_from_domain,
|
|
register_public_ip_in_contact: c.register_public_ip_in_contact,
|
|
outbound_sip_proxy: c.outbound_sip_proxy,
|
|
account_sid: c.account_sid,
|
|
trunk_type: c.trunk_type || 'static_ip'
|
|
};
|
|
}
|
|
|
|
async function getLocalSIPDomain(logger, srf) {
|
|
const { lookupSystemInformation } = srf.locals.dbHelpers;
|
|
try {
|
|
const systemInfo = await lookupSystemInformation();
|
|
if (systemInfo) {
|
|
logger.info(`lookup of sip domain from system_information: ${systemInfo.sip_domain_name}`);
|
|
srf.locals.localSIPDomain = systemInfo.sip_domain_name;
|
|
}
|
|
else {
|
|
logger.info('no system_information found, we will use the realm or public ip as the domain');
|
|
return false;
|
|
}
|
|
} catch (err) {
|
|
logger.info({ err }, 'Error looking up system information');
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Filters gateway array to remove duplicates based on ipv4, sip_realm, username and password
|
|
* @param {Array} gateways - Array of gateway objects
|
|
* @param {Object} logger - Logger instance to log duplicate entries
|
|
* @returns {Array} - Filtered array with unique gateways
|
|
*/
|
|
function getUniqueGateways(gateways, logger) {
|
|
const uniqueGatewayKeys = new Set();
|
|
const duplicateCounts = new Map();
|
|
const uniqueGateways = [];
|
|
|
|
for (const gw of gateways) {
|
|
const key = `${gw.ipv4}:${gw.sip_realm}:${gw.carrier?.register_username}:${gw.carrier?.register_password}`;
|
|
if (!gw.carrier?.register_password) {
|
|
logger.info({gw}, `Gateway ${key} does not have a password, ignoring`);
|
|
continue;
|
|
}
|
|
|
|
if (uniqueGatewayKeys.has(key)) {
|
|
duplicateCounts.set(key, (duplicateCounts.get(key) || 1) + 1);
|
|
} else {
|
|
uniqueGatewayKeys.add(key);
|
|
uniqueGateways.push(gw);
|
|
}
|
|
}
|
|
|
|
// Log summary for each duplicate gateway
|
|
for (const [key, count] of duplicateCounts) {
|
|
logger.info({key, count}, `Found ${count} duplicate gateways for ${key}, ignoring duplicates`);
|
|
}
|
|
|
|
return uniqueGateways;
|
|
}
|
|
|
|
module.exports = async(logger, srf) => {
|
|
if (initialized) return;
|
|
initialized = true;
|
|
const { addKeyNx } = srf.locals.realtimeDbHelpers;
|
|
const myToken = short.generate();
|
|
srf.locals.regbot = {
|
|
myToken,
|
|
active: false
|
|
};
|
|
|
|
srf.locals.regbotStatus = () => {
|
|
return {
|
|
total: regbots.length,
|
|
registered: getCountSuccessfulRegbots(),
|
|
active: srf.locals.regbot.active
|
|
};
|
|
};
|
|
|
|
/* Set the Local SIP domain on srf.locals */
|
|
await getLocalSIPDomain(logger, srf); // Initial Setup
|
|
setInterval(getLocalSIPDomain, 300000, logger, srf); //Refresh SIP Domain every 5 mins
|
|
|
|
/* sleep a random duration between 0 and MAX_INITIAL_DELAY seconds */
|
|
const ms = Math.floor(Math.random() * MAX_INITIAL_DELAY) * 1000;
|
|
logger.info(`waiting ${ms}ms before attempting to claim regbot responsibility with token ${myToken}`);
|
|
await waitFor(ms);
|
|
|
|
/* try to claim responsibility */
|
|
const result = await addKeyNx(regbotKey, myToken, REGBOT_STATUS_CHECK_INTERVAL + 10);
|
|
if (result === 'OK') {
|
|
srf.locals.regbot.active = true;
|
|
logger.info(`successfully claimed regbot responsibility with token ${myToken}`);
|
|
}
|
|
else {
|
|
logger.info(`failed to claim regbot responsibility with my token ${myToken}`);
|
|
}
|
|
|
|
/* check every so often if I need to go from inactive->active (or vice versa) */
|
|
setInterval(checkStatus.bind(null, logger, srf), REGBOT_STATUS_CHECK_INTERVAL * 1000);
|
|
|
|
/* on an AWS scale-in, hand the regbot role to another SBC while we are still up */
|
|
if (AWS_LIFECYCLE_DRAIN) {
|
|
require('./aws-lifecycle')(logger, () => {
|
|
handoff(logger, srf).catch((err) => logger.error({err}, 'regbot handoff failed'));
|
|
});
|
|
}
|
|
|
|
/* if I am the regbot holder, then kick it off */
|
|
if (srf.locals.regbot.active) {
|
|
updateCarrierRegbots(logger, srf)
|
|
.catch((err) => {
|
|
logger.error({ err }, 'updateCarrierRegbots failure');
|
|
});
|
|
}
|
|
|
|
return srf.locals.regbot.active;
|
|
};
|
|
|
|
const checkStatus = async(logger, srf) => {
|
|
const { addKeyNx, addKey, retrieveKey } = srf.locals.realtimeDbHelpers;
|
|
const { myToken, active, draining } = srf.locals.regbot;
|
|
|
|
/* a draining SBC has handed off the role and must never take it back */
|
|
if (draining) return;
|
|
|
|
logger.info({ active, myToken }, 'checking in on regbot status');
|
|
try {
|
|
const token = await retrieveKey(regbotKey);
|
|
let grabForTheWheel = false;
|
|
|
|
if (active) {
|
|
if (token === myToken) {
|
|
logger.info('I am active, and shall continue in my role as regbot');
|
|
addKey(regbotKey, myToken, REGBOT_STATUS_CHECK_INTERVAL + 10)
|
|
.then(updateCarrierRegbots.bind(null, logger, srf))
|
|
.catch((err) => {
|
|
logger.error({ err }, 'updateCarrierRegbots failure');
|
|
});
|
|
}
|
|
else if (token && token !== myToken) {
|
|
logger.info('Someone else grabbed the role! I need to stand down');
|
|
srf.locals.regbot.active = false;
|
|
regbots.forEach((rb) => rb.stop(srf));
|
|
regbots.length = 0;
|
|
/* reset hashes so the next updateCarrierRegbots rebuilds from scratch;
|
|
otherwise the array stays empty until a carrier config change */
|
|
carriersHash = '';
|
|
gatewaysHash = '';
|
|
}
|
|
else {
|
|
grabForTheWheel = true;
|
|
regbots.forEach((rb) => rb.stop(srf));
|
|
regbots.length = 0;
|
|
/* reset hashes so the next updateCarrierRegbots rebuilds from scratch;
|
|
otherwise the array stays empty until a carrier config change */
|
|
carriersHash = '';
|
|
gatewaysHash = '';
|
|
}
|
|
}
|
|
else {
|
|
if (token) {
|
|
logger.info('I am inactive and someone else is performing the role');
|
|
}
|
|
else {
|
|
grabForTheWheel = true;
|
|
}
|
|
}
|
|
|
|
if (grabForTheWheel) {
|
|
logger.info('regbot status is vacated, try to grab it!');
|
|
const result = await addKeyNx(regbotKey, myToken, REGBOT_STATUS_CHECK_INTERVAL + 10);
|
|
if (result === 'OK') {
|
|
srf.locals.regbot.active = true;
|
|
logger.info(`successfully claimed regbot responsibility with token ${myToken}`);
|
|
updateCarrierRegbots(logger, srf)
|
|
.catch((err) => {
|
|
logger.error({ err }, 'updateCarrierRegbots failure');
|
|
});
|
|
}
|
|
else {
|
|
srf.locals.regbot.active = false;
|
|
logger.info('failed to claim regbot responsibility');
|
|
}
|
|
}
|
|
} catch (err) {
|
|
logger.error({ err }, 'checkStatus: ERROR');
|
|
}
|
|
};
|
|
|
|
const updateCarrierRegbots = async(logger, srf) => {
|
|
if (rebuildInProgress) {
|
|
logger.info('updateCarrierRegbots: rebuild already in progress, skipping');
|
|
return;
|
|
}
|
|
rebuildInProgress = true;
|
|
const { lookupAllVoipCarriers, lookupSipGatewaysByCarrier, lookupAccountBySid } = srf.locals.dbHelpers;
|
|
try {
|
|
|
|
/* first check: has anything changed (new carriers or gateways)? */
|
|
let hasChanged = false;
|
|
const gws = [];
|
|
const cs = (await lookupAllVoipCarriers())
|
|
.filter((c) => c.requires_register && c.is_active)
|
|
.map((c) => pickRelevantCarrierProperties(c));
|
|
const newCarriersHash = computeHash(cs);
|
|
if (newCarriersHash !== carriersHash) hasChanged = true;
|
|
for (const c of cs) {
|
|
try {
|
|
const arr = (await lookupSipGatewaysByCarrier(c.voip_carrier_sid))
|
|
.filter((gw) => gw.outbound && gw.is_active)
|
|
.map((gw) => {
|
|
gw.carrier = pickRelevantCarrierProperties(c);
|
|
return gw;
|
|
});
|
|
gws.push(...arr);
|
|
} catch (err) {
|
|
logger.error({ err }, 'updateCarrierRegbots Error retrieving gateways');
|
|
}
|
|
}
|
|
const newGatewaysHash = computeHash(gws);
|
|
if (newGatewaysHash !== gatewaysHash) hasChanged = true;
|
|
if (hasChanged) {
|
|
|
|
debug('updateCarrierRegbots: got new or changed carriers');
|
|
logger.info({count: gws.length}, 'updateCarrierRegbots: got new or changed carriers');
|
|
|
|
carriersHash = newCarriersHash;
|
|
gatewaysHash = newGatewaysHash;
|
|
|
|
// Build maps of existing regbots for O(1) lookup
|
|
const existingByKey = new Map();
|
|
const existingByCarrierIpPort = new Map();
|
|
for (const rb of regbots) {
|
|
existingByKey.set(rb.configKey(), rb);
|
|
existingByCarrierIpPort.set(`${rb.voip_carrier_sid}:${rb.ipv4}:${rb.port}`, rb);
|
|
}
|
|
|
|
const newRegbots = [];
|
|
const newCarrierSids = new Set();
|
|
const keepKeys = new Set();
|
|
const accountSipRealmCache = new Map();
|
|
let batch_count = 0;
|
|
|
|
for (const gw of getUniqueGateways(gws, logger)) {
|
|
let accountSipRealm;
|
|
if (!gw.carrier.register_public_ip_in_contact && gw.carrier.account_sid) {
|
|
const acctSid = gw.carrier.account_sid;
|
|
if (accountSipRealmCache.has(acctSid)) {
|
|
accountSipRealm = accountSipRealmCache.get(acctSid);
|
|
} else {
|
|
const account = await lookupAccountBySid(acctSid);
|
|
accountSipRealm = (account && account.sip_realm) || null;
|
|
accountSipRealmCache.set(acctSid, accountSipRealm);
|
|
}
|
|
}
|
|
try {
|
|
const opts = {
|
|
voip_carrier_sid: gw.carrier.voip_carrier_sid,
|
|
account_sip_realm: accountSipRealm,
|
|
ipv4: gw.ipv4,
|
|
port: gw.port,
|
|
protocol: gw.protocol,
|
|
use_sips_scheme: gw.use_sips_scheme,
|
|
username: gw.carrier.register_username,
|
|
password: gw.carrier.register_password,
|
|
sip_realm: gw.carrier.register_sip_realm,
|
|
from_user: gw.carrier.register_from_user,
|
|
from_domain: gw.carrier.register_from_domain,
|
|
use_public_ip_in_contact: gw.carrier.register_public_ip_in_contact,
|
|
outbound_sip_proxy: gw.carrier.outbound_sip_proxy,
|
|
trunk_type: gw.carrier.trunk_type,
|
|
sip_gateway_sid: gw.sip_gateway_sid
|
|
};
|
|
|
|
const key = Regbot.configKeyFromOpts(opts);
|
|
|
|
if (existingByKey.has(key)) {
|
|
// Unchanged regbot -- keep the existing one running
|
|
const existing = existingByKey.get(key);
|
|
newRegbots.push(existing);
|
|
newCarrierSids.add(existing.voip_carrier_sid);
|
|
keepKeys.add(key);
|
|
} else {
|
|
// New or changed regbot -- only now construct the instance
|
|
const rb = new Regbot(logger, opts);
|
|
const oldRb = existingByCarrierIpPort.get(
|
|
`${rb.voip_carrier_sid}:${rb.ipv4}:${rb.port}`);
|
|
if (oldRb) {
|
|
rb.consecutiveRemoveFailures = oldRb.consecutiveRemoveFailures || 0;
|
|
}
|
|
newRegbots.push(rb);
|
|
newCarrierSids.add(rb.voip_carrier_sid);
|
|
rb.start(srf);
|
|
batch_count++;
|
|
if (batch_count >= JAMBONES_REGBOT_BATCH_SIZE) {
|
|
batch_count = 0;
|
|
await sleepFor(JAMBONES_REGBOT_BATCH_SLEEP_MS);
|
|
}
|
|
}
|
|
} catch (err) {
|
|
const { updateVoipCarriersRegisterStatus } = srf.locals.dbHelpers;
|
|
updateVoipCarriersRegisterStatus(gw.carrier.voip_carrier_sid, JSON.stringify({
|
|
status: 'fail',
|
|
reason: err.message,
|
|
}));
|
|
logger.error({ err },
|
|
`Error starting regbot for ${gw.carrier.register_username}@${gw.carrier.register_sip_realm}`);
|
|
}
|
|
}
|
|
|
|
// Stop old regbots that are no longer needed
|
|
for (const [key, rb] of existingByKey) {
|
|
if (!keepKeys.has(key)) {
|
|
// Check if a replacement regbot exists for the same carrier
|
|
const hasReplacement = newCarrierSids.has(rb.voip_carrier_sid);
|
|
if (hasReplacement) {
|
|
// Config changed but carrier still active -- stop timer only,
|
|
// keep ephemeral gateways in Redis until new regbot overwrites them
|
|
rb.stopTimer();
|
|
logger.info(`config changed for regbot ${rb.aor}, preserving gateways until re-registered`);
|
|
} else {
|
|
// Carrier removed or deactivated -- full cleanup
|
|
rb.stop(srf);
|
|
logger.info(`removed regbot ${rb.aor}, deleted ephemeral gateways`);
|
|
}
|
|
}
|
|
}
|
|
|
|
// Replace the regbots array (chunked to avoid call stack argument limit)
|
|
regbots.length = 0;
|
|
for (let i = 0; i < newRegbots.length; i += 1000) {
|
|
Array.prototype.push.apply(regbots, newRegbots.slice(i, i + 1000));
|
|
}
|
|
|
|
logger.info(`updateCarrierRegbots: ${regbots.length} regbots active, ` +
|
|
`${keepKeys.size} kept, ${regbots.length - keepKeys.size} new`);
|
|
}
|
|
} catch (err) {
|
|
logger.error({ err }, 'updateCarrierRegbots Error');
|
|
} finally {
|
|
rebuildInProgress = false;
|
|
}
|
|
};
|
|
|
|
/**
|
|
* Called when drachtio reconnects after a server restart.
|
|
* The regbots survive in memory but their SIP state on the (restarted) drachtio server is gone,
|
|
* so force each one to re-register immediately. If the array was emptied (e.g. by checkStatus)
|
|
* but this instance is still the active holder, force a full rebuild instead.
|
|
*/
|
|
const resync = async(logger, srf) => {
|
|
if (!srf.locals.regbot) return;
|
|
if (!srf.locals.regbot.active) {
|
|
logger.info('drachtio reconnect: not the active regbot holder, nothing to re-register');
|
|
return;
|
|
}
|
|
if (regbots.length === 0) {
|
|
/* nothing primed (e.g. cleared by checkStatus) -- force a full rebuild */
|
|
carriersHash = '';
|
|
gatewaysHash = '';
|
|
return updateCarrierRegbots(logger, srf);
|
|
}
|
|
/* A registered bot with a live refresh timer keeps its binding at the carrier across a
|
|
drachtio restart -- REGISTER is transaction-stateless on drachtio -- so leave it alone.
|
|
Only bots with a REGISTER in flight (watchdog set) or not currently registered need to be
|
|
re-driven; batch them so a large fleet does not REGISTER-storm all at once. */
|
|
const broken = regbots.filter((rb) => rb.status !== 'registered' || rb.watchdog);
|
|
if (broken.length === 0) {
|
|
logger.info(`drachtio reconnect: all ${regbots.length} regbots healthy, nothing to re-register`);
|
|
return;
|
|
}
|
|
logger.info(`drachtio reconnect: re-registering ${broken.length} of ${regbots.length} regbots`);
|
|
let batch_count = 0;
|
|
for (const rb of broken) {
|
|
rb.reregister(srf);
|
|
if (++batch_count >= JAMBONES_REGBOT_BATCH_SIZE) {
|
|
batch_count = 0;
|
|
await sleepFor(JAMBONES_REGBOT_BATCH_SLEEP_MS);
|
|
}
|
|
}
|
|
};
|
|
|
|
/**
|
|
* Called when this SBC starts draining for scale-in. If we hold the regbot role, give it up while
|
|
* we are still running rather than letting the lease expire after we are gone: the draining SBC
|
|
* answers new INVITEs with 503, and the carriers would otherwise keep sending calls for
|
|
* registration trunks here until the other SBC noticed the lapsed lease.
|
|
*
|
|
* Once another SBC has claimed the role (and so is registering with its own Call-ID), remove the
|
|
* bindings whose Contact points at this SBC, so registrars stop routing to an address that is
|
|
* going away.
|
|
*/
|
|
const handoff = async(logger, srf, opts = {}) => {
|
|
const {
|
|
pollMs = HANDOFF_POLL_MS,
|
|
timeoutMs = HANDOFF_TIMEOUT_MS,
|
|
graceMs = HANDOFF_GRACE_MS
|
|
} = opts;
|
|
const { retrieveKey, deleteKey } = srf.locals.realtimeDbHelpers;
|
|
const { myToken } = srf.locals.regbot;
|
|
|
|
srf.locals.regbot.draining = true;
|
|
if (!srf.locals.regbot.active) {
|
|
logger.info('scale-in: not the regbot holder, nothing to hand off');
|
|
return;
|
|
}
|
|
|
|
/* let an in-flight rebuild finish so it cannot start regbots after we stop them */
|
|
while (rebuildInProgress) await waitFor(250);
|
|
|
|
srf.locals.regbot.active = false;
|
|
const handedOff = regbots.splice(0, regbots.length);
|
|
/* stop refreshing but keep the ephemeral gateways: the new holder overwrites them */
|
|
handedOff.forEach((rb) => rb.stopTimer());
|
|
carriersHash = '';
|
|
gatewaysHash = '';
|
|
|
|
if (await retrieveKey(regbotKey) === myToken) await deleteKey(regbotKey);
|
|
logger.info(`scale-in: released regbot role (${handedOff.length} regbots), waiting for another SBC to claim it`);
|
|
|
|
let successor;
|
|
for (const deadline = Date.now() + timeoutMs; Date.now() < deadline;) {
|
|
await waitFor(pollMs);
|
|
const token = await retrieveKey(regbotKey);
|
|
if (token && token !== myToken) {
|
|
successor = token;
|
|
break;
|
|
}
|
|
}
|
|
if (!successor) {
|
|
logger.info('scale-in: no other SBC claimed the regbot role; leaving our bindings to expire');
|
|
return;
|
|
}
|
|
|
|
/* the successor registers in batches from the moment it claims the role; trail it by graceMs
|
|
at the same pace so each of its bindings is in place before we remove the matching old one */
|
|
await waitFor(graceMs);
|
|
const ours = handedOff.filter((rb) => rb.use_public_ip_in_contact);
|
|
logger.info(`scale-in: SBC ${successor} now holds the regbot role; ` +
|
|
`un-registering ${ours.length} bindings that point at this SBC`);
|
|
let batch_count = 0;
|
|
for (const rb of ours) {
|
|
await rb.unregister(srf);
|
|
if (++batch_count >= JAMBONES_REGBOT_BATCH_SIZE) {
|
|
batch_count = 0;
|
|
await sleepFor(JAMBONES_REGBOT_BATCH_SLEEP_MS);
|
|
}
|
|
}
|
|
};
|
|
|
|
module.exports.resync = resync;
|
|
module.exports.handoff = handoff;
|
|
|
|
// exposed for unit testing: stop all regbots and reset module state
|
|
module.exports._addForTest = (rb) => regbots.push(rb);
|
|
module.exports._resetForTest = () => {
|
|
regbots.forEach((rb) => {
|
|
rb.retired = true;
|
|
clearTimeout(rb.timer);
|
|
clearTimeout(rb.watchdog);
|
|
});
|
|
regbots.length = 0;
|
|
carriersHash = '';
|
|
gatewaysHash = '';
|
|
};
|