mirror of
https://github.com/jambonz/sbc-sip-sidecar.git
synced 2026-10-04 02:04:20 +00:00
* restart regbots on drachtio reconnect * PR feedback * add check for multiple regbot instances
456 lines
17 KiB
JavaScript
456 lines
17 KiB
JavaScript
const dns = require('dns').promises;
|
|
const {
|
|
JAMBONES_REGBOT_DEFAULT_EXPIRES_INTERVAL,
|
|
JAMBONES_REGBOT_MIN_EXPIRES_INTERVAL,
|
|
JAMBONES_REGBOT_CONTACT_USE_IP,
|
|
REGISTER_RESPONSE_REMOVE,
|
|
JAMBONES_REGBOT_USER_AGENT,
|
|
JAMBONES_REGBOT_FAILURE_RETRY_INTERVAL,
|
|
JAMBONES_REGBOT_REGISTER_FAILURE_THRESHOLD,
|
|
JAMBONES_REGBOT_RESPONSE_TIMEOUT
|
|
} = require('./config');
|
|
const {isValidDomainOrIP, isValidIPv4} = require('./utils');
|
|
const parseUri = require('drachtio-srf').parseUri;
|
|
const DEFAULT_EXPIRES = (parseInt(JAMBONES_REGBOT_DEFAULT_EXPIRES_INTERVAL) || 3600);
|
|
const MIN_EXPIRES = (parseInt(JAMBONES_REGBOT_MIN_EXPIRES_INTERVAL) || 90);
|
|
const FAILURE_RETRY_INTERVAL = (parseInt(JAMBONES_REGBOT_FAILURE_RETRY_INTERVAL) || 300);
|
|
const REGISTER_FAILURE_THRESHOLD = (parseInt(JAMBONES_REGBOT_REGISTER_FAILURE_THRESHOLD) || 3);
|
|
// Backstop for a REGISTER whose response never arrives at all (dropped socket). Kept above
|
|
// drachtio's own non-INVITE Timer F (~32s) so its 408 normally reaches us first and drives the
|
|
// usual fail/backoff; this only fires when even that never comes.
|
|
const RESPONSE_TIMEOUT = (parseInt(JAMBONES_REGBOT_RESPONSE_TIMEOUT) || 45000);
|
|
const assert = require('assert');
|
|
const version = require('../package.json').version;
|
|
const useragent = JAMBONES_REGBOT_USER_AGENT || `Jambonz ${version}`;
|
|
|
|
/**
|
|
* A 200 OK to a REGISTER may echo back every Contact bound to the AoR, not just ours
|
|
* (e.g. when another device is registered against the same account). Per RFC 3261 10.2.4
|
|
* we must refresh based on the expires of *our own* binding, so locate the returned Contact
|
|
* whose user and host match the Contact we sent.
|
|
* Returns the parsed expires (a number) of our binding, or undefined if it can't be matched.
|
|
*/
|
|
const getOwnContactExpires = (contacts, ownContactUri) => {
|
|
let own;
|
|
try {
|
|
own = parseUri(ownContactUri);
|
|
} catch {
|
|
return undefined;
|
|
}
|
|
if (!own || !own.host) return undefined;
|
|
const ownUser = own.user;
|
|
const ownHost = own.host.toLowerCase();
|
|
for (const c of contacts) {
|
|
let u;
|
|
try {
|
|
u = parseUri(c.uri);
|
|
} catch {
|
|
continue;
|
|
}
|
|
if (u && u.user === ownUser && (u.host || '').toLowerCase() === ownHost &&
|
|
c.params && c.params.expires !== undefined) {
|
|
return parseInt(c.params.expires);
|
|
}
|
|
}
|
|
return undefined;
|
|
};
|
|
|
|
class Regbot {
|
|
constructor(logger, opts) {
|
|
this.logger = logger;
|
|
|
|
[
|
|
'voip_carrier_sid',
|
|
'ipv4',
|
|
'port',
|
|
'username',
|
|
'password',
|
|
'protocol',
|
|
'account_sip_realm',
|
|
'outbound_sip_proxy',
|
|
'trunk_type',
|
|
'sip_gateway_sid'
|
|
].forEach((prop) => this[prop] = opts[prop]);
|
|
|
|
this.sip_realm = opts.sip_realm || opts.ipv4;
|
|
this.use_public_ip_in_contact = opts.use_public_ip_in_contact || JAMBONES_REGBOT_CONTACT_USE_IP;
|
|
this.use_sips_scheme = opts.use_sips_scheme || false;
|
|
|
|
this.fromUser = opts.from_user || this.username;
|
|
const fromDomain = opts.from_domain || this.sip_realm;
|
|
if (!isValidDomainOrIP(fromDomain)) {
|
|
throw new Error(`Invalid from_domain ${fromDomain}`);
|
|
}
|
|
this.from = `sip:${this.fromUser}@${fromDomain}`;
|
|
this.aor = `${this.fromUser}@${this.sip_realm}`;
|
|
this.status = 'none';
|
|
this.consecutiveRemoveFailures = 0;
|
|
// monotonic per-attempt counter; only the latest attempt's response/watchdog is acted upon
|
|
this.epoch = 0;
|
|
}
|
|
|
|
async start(srf) {
|
|
assert(!this.timer);
|
|
|
|
this.logger.info(`starting regbot for ${this.fromUser}@${this.sip_realm}`);
|
|
this.register(srf);
|
|
}
|
|
|
|
/**
|
|
* Force an immediate re-registration, e.g. after the drachtio server restarts.
|
|
* Clears any pending (possibly stale) re-register timer and sends a fresh REGISTER now;
|
|
* register() re-arms the timer chain from the response.
|
|
*/
|
|
reregister(srf) {
|
|
if (this.retired) return;
|
|
clearTimeout(this.timer);
|
|
this.timer = null;
|
|
this.register(srf);
|
|
}
|
|
|
|
stop(srf) {
|
|
this.retired = true;
|
|
const { deleteEphemeralGateway } = srf.locals.realtimeDbHelpers;
|
|
clearTimeout(this.timer);
|
|
this.timer = null;
|
|
clearTimeout(this.watchdog);
|
|
this.watchdog = null;
|
|
// remove any ephemeral gateways created for this regbot
|
|
if (this.addresses && this.addresses.length) {
|
|
this.addresses.forEach((ip) => {
|
|
deleteEphemeralGateway(ip, this.voip_carrier_sid).catch((err) => {
|
|
this.logger.error({err, ip}, 'Error deleting ephemeral gateway on regbot stop');
|
|
});
|
|
});
|
|
}
|
|
|
|
}
|
|
|
|
stopTimer() {
|
|
this.retired = true;
|
|
clearTimeout(this.timer);
|
|
this.timer = null;
|
|
clearTimeout(this.watchdog);
|
|
this.watchdog = null;
|
|
}
|
|
|
|
configKey() {
|
|
return [
|
|
this.voip_carrier_sid, this.ipv4, this.port,
|
|
this.username, this.password, this.sip_realm,
|
|
this.protocol, this.use_sips_scheme,
|
|
this.use_public_ip_in_contact, this.outbound_sip_proxy,
|
|
this.trunk_type, this.sip_gateway_sid,
|
|
this.account_sip_realm, this.fromUser, this.from
|
|
].join('|');
|
|
}
|
|
|
|
static configKeyFromOpts(opts) {
|
|
const sip_realm = opts.sip_realm || opts.ipv4;
|
|
const fromUser = opts.from_user || opts.username;
|
|
const fromDomain = opts.from_domain || sip_realm;
|
|
return [
|
|
opts.voip_carrier_sid, opts.ipv4, opts.port,
|
|
opts.username, opts.password, sip_realm,
|
|
opts.protocol, opts.use_sips_scheme || false,
|
|
opts.use_public_ip_in_contact || JAMBONES_REGBOT_CONTACT_USE_IP,
|
|
opts.outbound_sip_proxy,
|
|
opts.trunk_type, opts.sip_gateway_sid,
|
|
opts.account_sip_realm, fromUser, `sip:${fromUser}@${fromDomain}`
|
|
].join('|');
|
|
}
|
|
|
|
toJSON() {
|
|
return {
|
|
voip_carrier_sid: this.voip_carrier_sid,
|
|
username: this.username,
|
|
fromUser: this.fromUser,
|
|
sip_realm: this.sip_realm,
|
|
ipv4: this.ipv4,
|
|
port: this.port,
|
|
aor: this.aor,
|
|
status: this.status
|
|
};
|
|
}
|
|
|
|
async register(srf) {
|
|
const { createEphemeralGateway } = srf.locals.realtimeDbHelpers;
|
|
const { updateVoipCarriersRegisterStatus } = srf.locals.dbHelpers;
|
|
const { writeAlerts, localSIPDomain } = srf.locals;
|
|
// stamp this attempt so a late/stale response (e.g. after a watchdog-driven retry) is ignored
|
|
const epoch = ++this.epoch;
|
|
try {
|
|
// transport
|
|
const transport = (this.protocol.includes('/') ? this.protocol.substring(0, this.protocol.indexOf('/')) :
|
|
this.protocol).toLowerCase();
|
|
|
|
// scheme
|
|
let scheme = 'sip';
|
|
if (transport === 'tls' && this.use_sips_scheme) scheme = 'sips';
|
|
|
|
let publicAddress = srf.locals.sbcPublicIpAddress.udp;
|
|
if (transport !== 'udp') {
|
|
if (srf.locals.sbcPublicIpAddress[transport]) {
|
|
publicAddress = srf.locals.sbcPublicIpAddress[transport];
|
|
}
|
|
else if (transport === 'tls') {
|
|
publicAddress = srf.locals.sbcPublicIpAddress.udp;
|
|
}
|
|
}
|
|
|
|
let contactAddress = this.aor;
|
|
if (this.use_public_ip_in_contact) {
|
|
contactAddress = `${this.fromUser}@${publicAddress}`;
|
|
}
|
|
else if (this.account_sip_realm) {
|
|
contactAddress = `${this.fromUser}@${this.account_sip_realm}`;
|
|
}
|
|
else if (localSIPDomain) {
|
|
contactAddress = `${this.fromUser}@${localSIPDomain}`;
|
|
}
|
|
|
|
this.logger.debug(`sending REGISTER for ${this.aor}`);
|
|
|
|
let proxy;
|
|
if (this.outbound_sip_proxy) {
|
|
proxy = `sip:${this.outbound_sip_proxy};transport=${transport}`;
|
|
this.logger.debug(`sending via proxy ${proxy}`);
|
|
} else {
|
|
const isIPv4 = isValidIPv4(this.ipv4);
|
|
proxy = `sip:${this.ipv4}${isIPv4 ? `:${this.port}` : ''};transport=${transport}`;
|
|
this.logger.debug(`sending to registrar ${proxy}`);
|
|
}
|
|
/* Safety net for a REGISTER whose SIP response never arrives at all (srf.request has no
|
|
timeout) -- e.g. drachtio dropped the socket mid-transaction so we never even get its
|
|
408. Treat it exactly like the failure path (mark fail, record status, back off on
|
|
FAILURE_RETRY_INTERVAL) rather than re-sending inline, so an unreachable registrar does
|
|
not loop every RESPONSE_TIMEOUT. Bump the epoch so a late response is ignored. */
|
|
clearTimeout(this.watchdog);
|
|
this.watchdog = setTimeout(() => {
|
|
if (this.retired || epoch !== this.epoch) return;
|
|
this.watchdog = null;
|
|
this.epoch++;
|
|
this.status = 'fail';
|
|
this.logger.info(`${this.aor}: no response to REGISTER within ${RESPONSE_TIMEOUT}ms, backing off`);
|
|
updateVoipCarriersRegisterStatus(this.voip_carrier_sid, JSON.stringify({
|
|
status: 'fail',
|
|
reason: `no response within ${RESPONSE_TIMEOUT}ms`
|
|
}));
|
|
this.timer = setTimeout(this.register.bind(this, srf), FAILURE_RETRY_INTERVAL * 1000);
|
|
}, RESPONSE_TIMEOUT);
|
|
const req = await srf.request(`${scheme}:${this.sip_realm}`, {
|
|
method: 'REGISTER',
|
|
proxy,
|
|
headers: {
|
|
'Call-ID': this.sip_gateway_sid,
|
|
'From': this.from,
|
|
'To': this.from,
|
|
'Contact': `<${scheme}:${contactAddress};transport=${transport}>;expires=${DEFAULT_EXPIRES}`,
|
|
'Expires': DEFAULT_EXPIRES,
|
|
'User-Agent': useragent
|
|
},
|
|
auth: {
|
|
username: this.username,
|
|
password: this.password
|
|
}
|
|
});
|
|
req.on('response', async(res) => {
|
|
if (epoch !== this.epoch) {
|
|
this.logger.info(`${this.aor}: ignoring stale REGISTER response (superseded by a newer attempt)`);
|
|
return;
|
|
}
|
|
clearTimeout(this.watchdog);
|
|
this.watchdog = null;
|
|
if (this.retired) {
|
|
this.logger.info(`${this.aor}: ignoring response, regbot has been retired`);
|
|
return;
|
|
}
|
|
let expires;
|
|
if (res.status !== 200) {
|
|
this.status = 'fail';
|
|
this.logger.info(`${this.aor}: got ${res.status} registering to ${this.ipv4}:${this.port}`);
|
|
this.timer = setTimeout(this.register.bind(this, srf), FAILURE_RETRY_INTERVAL * 1000);
|
|
if (REGISTER_RESPONSE_REMOVE.includes(res.status)) {
|
|
this.consecutiveRemoveFailures++;
|
|
// eslint-disable-next-line max-len
|
|
this.logger.info(`${this.aor}: consecutive remove-failures: ${this.consecutiveRemoveFailures}/${REGISTER_FAILURE_THRESHOLD}`);
|
|
if (this.consecutiveRemoveFailures >= REGISTER_FAILURE_THRESHOLD) {
|
|
const { updateCarrierBySid, lookupCarrierBySid } = srf.locals.dbHelpers;
|
|
this.stop(srf); //Remove the retry timer
|
|
const carrier = await lookupCarrierBySid(this.voip_carrier_sid);
|
|
if (carrier && carrier.trunk_type != 'reg') {
|
|
await updateCarrierBySid(this.voip_carrier_sid, {requires_register: false});
|
|
// eslint-disable-next-line max-len
|
|
this.logger.info(`Disabling Outbound Registration for carrier ${carrier.name} (sid:${carrier.voip_carrier_sid})`);
|
|
writeAlerts({
|
|
account_sid: carrier.account_sid,
|
|
service_provider_sid: carrier.service_provider_sid,
|
|
// eslint-disable-next-line max-len
|
|
message: `Disabling Outbound Registration for carrier ${carrier.name} (sid:${carrier.voip_carrier_sid})`
|
|
});
|
|
} else {
|
|
await updateCarrierBySid(this.voip_carrier_sid, {is_active: false});
|
|
// eslint-disable-next-line max-len
|
|
this.logger.info(`Deactivating carrier ${carrier.name} (sid:${carrier.voip_carrier_sid}) due to registration errors`);
|
|
writeAlerts({
|
|
account_sid: carrier.account_sid,
|
|
service_provider_sid: carrier.service_provider_sid,
|
|
// eslint-disable-next-line max-len
|
|
message: `Deactivating carrier ${carrier.name} (sid:${carrier.voip_carrier_sid}) due to registration errors`
|
|
});
|
|
}
|
|
}
|
|
}
|
|
expires = 0;
|
|
}
|
|
else {
|
|
// the code parses the SIP headers to get the expires value
|
|
// if there is a Contact header, it will use the expires value from there
|
|
// otherwise, it will use the Expires header, acording to the SIP RFC 3261, section 10.2.4 Refreshing Bindings
|
|
this.status = 'registered';
|
|
this.consecutiveRemoveFailures = 0;
|
|
expires = DEFAULT_EXPIRES;
|
|
|
|
if (res.has('Expires')) {
|
|
expires = parseInt(res.get('Expires'));
|
|
this.logger.debug(`Using Expires header value of ${expires}`);
|
|
}
|
|
|
|
if (res.has('Contact')) {
|
|
const contacts = res.getParsedHeader('Contact');
|
|
const ownExpires = getOwnContactExpires(contacts, `${scheme}:${contactAddress}`);
|
|
if (typeof ownExpires === 'number' && !isNaN(ownExpires)) {
|
|
// refresh based on our own binding's expires
|
|
expires = ownExpires;
|
|
} else if (contacts.length === 1 && contacts[0].params && contacts[0].params.expires) {
|
|
// single binding returned: safe to use even if the registrar rewrote the contact host
|
|
expires = parseInt(contacts[0].params.expires);
|
|
} else if (contacts.length > 1) {
|
|
this.logger.info({aor: this.aor, ipv4: this.ipv4, port: this.port},
|
|
// eslint-disable-next-line max-len
|
|
`200 OK returned ${contacts.length} contacts but none matched our own binding; using expires ${expires}`);
|
|
}
|
|
} else {
|
|
this.logger.debug({ aor: this.aor, ipv4: this.ipv4, port: this.port },
|
|
'no Contact header in 200 OK');
|
|
}
|
|
|
|
if (isNaN(expires) || expires < MIN_EXPIRES) {
|
|
this.logger.debug({ aor: this.aor, ipv4: this.ipv4, port: this.port },
|
|
`got expires of ${expires} in 200 OK, too small so setting to ${MIN_EXPIRES}`);
|
|
expires = MIN_EXPIRES;
|
|
}
|
|
this.logger.debug(`setting timer for next register to ${expires} seconds`);
|
|
this.timer = setTimeout(this.register.bind(this, srf), (expires / 2) * 1000);
|
|
}
|
|
const timestamp = new Date().toISOString();
|
|
|
|
//update registration status for the carrier in the database
|
|
updateVoipCarriersRegisterStatus(this.voip_carrier_sid, JSON.stringify({
|
|
status: res.status === 200 ? 'ok' : 'fail',
|
|
reason: `${res.status} ${res.reason}`,
|
|
cseq: req.get('Cseq'),
|
|
callId: req.get('Call-Id'),
|
|
timestamp: timestamp,
|
|
expires: expires
|
|
}));
|
|
|
|
// for reg trunks, create ephemeral set of IP addresses for inbound gateways
|
|
if (this.trunk_type === 'reg') {
|
|
this.addresses = [];
|
|
if (isValidIPv4(this.ipv4)) {
|
|
this.addresses.push(this.ipv4);
|
|
}
|
|
else {
|
|
if (this.port) {
|
|
const addrs = await dnsResolverA(this.logger, this.ipv4);
|
|
this.addresses.push(...addrs);
|
|
}
|
|
else {
|
|
const addrs = await dnsResolverSrv(this.logger, this.ipv4, this.transport);
|
|
if (addrs.length) {
|
|
this.addresses.push(...addrs);
|
|
} else {
|
|
this.logger.info({ipv4: this.ipv4, transport: this.transport},
|
|
'No SRV addresses found for reg-gateway');
|
|
const addrsARecord = await dnsResolverA(this.logger, this.ipv4);
|
|
if (addrsARecord.length) {
|
|
this.addresses.push(...addrsARecord);
|
|
} else {
|
|
this.logger.info({ipv4: this.ipv4}, 'No A record found for reg-gateway');
|
|
}
|
|
}
|
|
}
|
|
}
|
|
if (this.addresses.length) {
|
|
try {
|
|
await Promise.all(
|
|
this.addresses.map((ip) => createEphemeralGateway(ip, this.voip_carrier_sid, expires))
|
|
);
|
|
} catch (err) {
|
|
this.logger.error({addresses: this.addresses, err}, 'Error creating hash for reg-gateway');
|
|
}
|
|
this.logger.debug({addresses: this.addresses},
|
|
`Created ephemeral gateways for registration trunk ${this.voip_carrier_sid}, ${this.sip_realm}`);
|
|
}
|
|
}
|
|
});
|
|
} catch (err) {
|
|
clearTimeout(this.watchdog);
|
|
this.watchdog = null;
|
|
this.logger.error({ err }, `${this.aor}: Error registering to ${this.ipv4}:${this.port}`);
|
|
this.timer = setTimeout(this.register.bind(this, srf), FAILURE_RETRY_INTERVAL * 1000);
|
|
updateVoipCarriersRegisterStatus(this.voip_carrier_sid, JSON.stringify({
|
|
status: 'fail',
|
|
reason: err
|
|
}));
|
|
}
|
|
|
|
}
|
|
}
|
|
|
|
const dnsResolverA = async(logger, hostname) => {
|
|
try {
|
|
const addresses = await dns.resolve4(hostname);
|
|
logger.debug({addresses}, `Regbot: resolved ${hostname} into ${addresses.length} IPs`);
|
|
return addresses;
|
|
} catch (err) {
|
|
logger.info({err}, `Error resolving ${hostname}`);
|
|
}
|
|
return [];
|
|
};
|
|
|
|
const dnsResolverSrv = async(logger, hostname, transport) => {
|
|
let name;
|
|
switch (transport) {
|
|
case 'tls':
|
|
name = `_sips._tcp.${hostname}`;
|
|
break;
|
|
case 'tcp':
|
|
name = `_sip._tcp.${hostname}`;
|
|
break;
|
|
default:
|
|
name = `_sip._udp.${hostname}`;
|
|
}
|
|
|
|
try {
|
|
const arr = await dns.resolveSrv(name);
|
|
logger.debug({arr}, `Regbot: resolved ${hostname}/${transport} into ${arr.length} results`);
|
|
const ips = await Promise.all(
|
|
arr.map((obj) => dnsResolverA(logger, obj.name))
|
|
);
|
|
return ips.flat();
|
|
}
|
|
catch (err) {
|
|
logger.info({err}, `SRV Error resolving ${hostname}`);
|
|
}
|
|
return [];
|
|
};
|
|
|
|
|
|
// exposed for unit testing
|
|
Regbot._getOwnContactExpires = getOwnContactExpires;
|
|
|
|
module.exports = Regbot;
|
|
|