From ce05d91e6d5b625fa1a75517d8879156ec4757fe Mon Sep 17 00:00:00 2001 From: Sam Machin Date: Wed, 9 Sep 2026 15:02:23 +0100 Subject: [PATCH] restart regbots on drachtio reconnect (#151) * restart regbots on drachtio reconnect * PR feedback * add check for multiple regbot instances --- app.js | 13 +++ lib/config.js | 3 + lib/regbot.js | 53 +++++++++++- lib/sip-trunk-register.js | 61 ++++++++++++++ test/index.js | 1 + test/regbot-reconnect-test.js | 130 +++++++++++++++++++++++++++++ test/regbot-unit-test.js | 150 ++++++++++++++++++++++++++++++++++ 7 files changed, 410 insertions(+), 1 deletion(-) create mode 100644 test/regbot-reconnect-test.js diff --git a/app.js b/app.js index 454395b..d5cdbf2 100644 --- a/app.js +++ b/app.js @@ -149,8 +149,12 @@ const matcher = new CIDRMatcher(cidrs); srf.locals.matcher = matcher; srf.connect({ host: DRACHTIO_HOST, port: DRACHTIO_PORT, secret: DRACHTIO_SECRET }); +let drachtioConnected = false; srf.on('connect', (err, hp, version, localHostports) => { if (err) return logger.error({ err }, 'Error connecting to drachtio server'); + // drachtio-srf re-emits 'connect' on every reconnect; distinguish a reconnect from first connect + const isReconnect = drachtioConnected; + drachtioConnected = true; logger.info(`connected to drachtio listening on ${hp}, local hostports: ${localHostports}`); if (localHostports) { @@ -244,6 +248,15 @@ srf.on('connect', (err, hp, version, localHostports) => { require('./lib/sip-trunk-register')(logger, srf); // Start Options bot require('./lib/sip-trunk-options-ping')(logger, srf); + + /* on a reconnect after a drachtio restart, the initial startup above is a no-op + (guarded by `initialized`), so force the surviving regbots to re-register */ + if (isReconnect) { + logger.info('drachtio reconnected — resyncing regbots'); + // eslint-disable-next-line promise/no-promise-in-callback + require('./lib/sip-trunk-register').resync(logger, srf) + .catch((e) => logger.error({ err: e }, 'regbot resync after reconnect failed')); + } }); if (NODE_ENV === 'test') { diff --git a/lib/config.js b/lib/config.js index 13f2ec6..c89df2b 100644 --- a/lib/config.js +++ b/lib/config.js @@ -44,6 +44,8 @@ const REGISTER_RESPONSE_REMOVE = process.env.REGISTER_RESPONSE_REMOVE?.split(',' const JAMBONES_REGBOT_USER_AGENT = process.env.JAMBONES_REGBOT_USER_AGENT ; const JAMBONES_REGBOT_FAILURE_RETRY_INTERVAL = process.env.JAMBONES_REGBOT_FAILURE_RETRY_INTERVAL; const JAMBONES_REGBOT_REGISTER_FAILURE_THRESHOLD = process.env.JAMBONES_REGBOT_REGISTER_FAILURE_THRESHOLD; +// how long (ms) to wait for a SIP response to a REGISTER before assuming it was lost and retrying +const JAMBONES_REGBOT_RESPONSE_TIMEOUT = process.env.JAMBONES_REGBOT_RESPONSE_TIMEOUT; /* Server control - external topology discovery and other server-control features (disabled unless truthy) */ const JAMBONES_SERVER_CONTROL = process.env.JAMBONES_SERVER_CONTROL; @@ -85,5 +87,6 @@ module.exports = { JAMBONES_REGBOT_USER_AGENT, JAMBONES_REGBOT_FAILURE_RETRY_INTERVAL, JAMBONES_REGBOT_REGISTER_FAILURE_THRESHOLD, + JAMBONES_REGBOT_RESPONSE_TIMEOUT, JAMBONES_SERVER_CONTROL }; diff --git a/lib/regbot.js b/lib/regbot.js index a21a9e1..90553f5 100644 --- a/lib/regbot.js +++ b/lib/regbot.js @@ -6,7 +6,8 @@ const { REGISTER_RESPONSE_REMOVE, JAMBONES_REGBOT_USER_AGENT, JAMBONES_REGBOT_FAILURE_RETRY_INTERVAL, - JAMBONES_REGBOT_REGISTER_FAILURE_THRESHOLD + JAMBONES_REGBOT_REGISTER_FAILURE_THRESHOLD, + JAMBONES_REGBOT_RESPONSE_TIMEOUT } = require('./config'); const {isValidDomainOrIP, isValidIPv4} = require('./utils'); const parseUri = require('drachtio-srf').parseUri; @@ -14,6 +15,10 @@ const DEFAULT_EXPIRES = (parseInt(JAMBONES_REGBOT_DEFAULT_EXPIRES_INTERVAL) || 3 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}`; @@ -80,6 +85,8 @@ class Regbot { 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) { @@ -89,11 +96,25 @@ class Regbot { 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) => { @@ -109,6 +130,8 @@ class Regbot { this.retired = true; clearTimeout(this.timer); this.timer = null; + clearTimeout(this.watchdog); + this.watchdog = null; } configKey() { @@ -154,6 +177,8 @@ class Regbot { 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('/')) : @@ -195,6 +220,24 @@ class Regbot { 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, @@ -212,6 +255,12 @@ class Regbot { } }); 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; @@ -347,6 +396,8 @@ class Regbot { } }); } 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({ diff --git a/lib/sip-trunk-register.js b/lib/sip-trunk-register.js index 827ba89..9069955 100644 --- a/lib/sip-trunk-register.js +++ b/lib/sip-trunk-register.js @@ -168,13 +168,22 @@ const checkStatus = async(logger, srf) => { } 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 { @@ -361,3 +370,55 @@ const updateCarrierRegbots = async(logger, srf) => { 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); + } + } +}; + +module.exports.resync = resync; + +// exposed for unit testing: stop all regbots and reset module state +module.exports._resetForTest = () => { + regbots.forEach((rb) => { + rb.retired = true; + clearTimeout(rb.timer); + clearTimeout(rb.watchdog); + }); + regbots.length = 0; + carriersHash = ''; + gatewaysHash = ''; +}; diff --git a/test/index.js b/test/index.js index f288ab9..019c19f 100644 --- a/test/index.js +++ b/test/index.js @@ -3,6 +3,7 @@ require('./create-test-db'); require('./regbot-tests'); require('./regbot-unit-test'); require('./regbot-concurrent-rebuild-test'); +require('./regbot-reconnect-test'); require('./sip-register-tests'); require('./sip-options-tests'); require('./cli-tests'); diff --git a/test/regbot-reconnect-test.js b/test/regbot-reconnect-test.js new file mode 100644 index 0000000..b82a2f0 --- /dev/null +++ b/test/regbot-reconnect-test.js @@ -0,0 +1,130 @@ +const test = require('tape'); +const { EventEmitter } = require('events'); +const clearModule = require('clear-module'); +const { + JAMBONES_LOGLEVEL, +} = require('../lib/config'); +const opts = Object.assign({ + timestamp: () => { return `, "time": "${new Date().toISOString()}"`; } +}, { level: JAMBONES_LOGLEVEL || 'info' }); +const logger = require('pino')(opts); + +/** + * Tests for drachtio-reconnect recovery (resync) in lib/sip-trunk-register.js. + * + * When the drachtio server restarts, the sidecar keeps running but the regbots' + * SIP state on the (restarted) server is gone. Previously nothing re-registered + * them: the reconnect path was guarded out by `initialized`, and + * updateCarrierRegbots only rebuilds on a config-hash change. resync() is the + * new entrypoint the reconnect handler calls. + */ + +const tick = () => new Promise((resolve) => setImmediate(resolve)); +const ok200 = { status: 200, reason: 'OK', has: () => false, get: () => undefined, getParsedHeader: () => [] }; + +// a mock srf that records each REGISTER and serves one carrier/gateway from the db-helpers +function makeSrf({ active }) { + const state = { requests: [] }; + const carriers = [{ + voip_carrier_sid: 'c1', + requires_register: 1, + is_active: 1, + register_username: 'user', + register_password: 'pass', + register_sip_realm: 'registrar.example.com', + register_public_ip_in_contact: 1, // avoids the account-realm lookup path + trunk_type: 'static_ip', // avoids the reg-trunk DNS/ephemeral-gateway path + account_sid: null + }]; + const gateways = [{ + sip_gateway_sid: 'gw1', + voip_carrier_sid: 'c1', + ipv4: '10.1.1.1', + port: 5060, + protocol: 'udp', + outbound: 1, + is_active: 1 + }]; + const srf = { + request: () => { + const req = new EventEmitter(); + req.get = () => ''; + state.requests.push(req); + return Promise.resolve(req); + }, + locals: { + regbot: { active, myToken: 'tok' }, + sbcPublicIpAddress: { udp: '203.0.113.1:5060' }, + localSIPDomain: 'sbc.example.com', + writeAlerts: () => {}, + realtimeDbHelpers: { + createEphemeralGateway: () => Promise.resolve(), + deleteEphemeralGateway: () => Promise.resolve() + }, + dbHelpers: { + updateVoipCarriersRegisterStatus: () => {}, + lookupAllVoipCarriers: () => Promise.resolve(carriers), + lookupSipGatewaysByCarrier: () => Promise.resolve(gateways), + lookupAccountBySid: () => Promise.resolve(null) + } + } + }; + return { srf, state }; +} + +test('resync is a no-op when this instance is not the active regbot holder', async(t) => { + clearModule('../lib/sip-trunk-register'); + const reg = require('../lib/sip-trunk-register'); + const { srf, state } = makeSrf({ active: false }); + + await reg.resync(logger, srf); + await tick(); + + t.equal(state.requests.length, 0, 'no REGISTER sent when not active'); + reg._resetForTest(); + t.end(); +}); + +test('resync rebuilds regbots when active but the array is empty (post-clear recovery)', async(t) => { + clearModule('../lib/sip-trunk-register'); + const reg = require('../lib/sip-trunk-register'); + const { srf, state } = makeSrf({ active: true }); + + // empty regbots + active holder -> resync resets hashes and forces a rebuild + await reg.resync(logger, srf); + await tick(); + await tick(); + t.equal(state.requests.length, 1, 'a regbot was rebuilt and sent a REGISTER'); + + reg._resetForTest(); + t.end(); +}); + +test('resync re-registers only broken regbots and leaves healthy ones alone', async(t) => { + clearModule('../lib/sip-trunk-register'); + const reg = require('../lib/sip-trunk-register'); + const { srf, state } = makeSrf({ active: true }); + + // build the bot; with no response yet it is unregistered (a REGISTER in flight) => broken + await reg.resync(logger, srf); + await tick(); + await tick(); + t.equal(state.requests.length, 1, 'a regbot was rebuilt and sent a REGISTER'); + + // reconnect while still unregistered -> it is re-driven + await reg.resync(logger, srf); + await tick(); + t.equal(state.requests.length, 2, 'broken (unregistered) regbot was re-registered on reconnect'); + + // let it register successfully + state.requests[1].emit('response', ok200); + await tick(); + + // reconnect again: a healthy registered bot keeps its binding, so it is skipped + await reg.resync(logger, srf); + await tick(); + t.equal(state.requests.length, 2, 'healthy registered regbot is not re-registered'); + + reg._resetForTest(); + t.end(); +}); diff --git a/test/regbot-unit-test.js b/test/regbot-unit-test.js index 9316745..1d6f920 100644 --- a/test/regbot-unit-test.js +++ b/test/regbot-unit-test.js @@ -1,4 +1,6 @@ const test = require('tape'); +const { EventEmitter } = require('events'); +const clearModule = require('clear-module'); const Regbot = require('../lib/regbot'); const { JAMBONES_LOGLEVEL, @@ -246,3 +248,151 @@ test('stop sets retired flag', (t) => { t.equal(rb.timer, null, 'timer is cleared'); t.end(); }); + +/* ---- re-registration / drachtio-reconnect recovery tests ---- */ + +const REGBOT_OPTS = { + voip_carrier_sid: 'carrier-1', + ipv4: '2.3.4.5', + port: 5060, + username: 'user', + password: 'password', + sip_realm: 'sip.server.com', + protocol: 'udp', + trunk_type: 'static_ip', // avoid the reg-trunk DNS/ephemeral-gateway path + sip_gateway_sid: 'gw-1' +}; + +// a mock srf whose request() records each REGISTER and returns a controllable fake req +function makeSrf() { + const state = { requests: [], statusUpdates: [] }; + const srf = { + request: (uri, opts) => { + const req = new EventEmitter(); + req.get = () => ''; + req.uri = uri; + req.opts = opts; + state.requests.push(req); + return Promise.resolve(req); + }, + locals: { + sbcPublicIpAddress: { udp: '203.0.113.1:5060' }, + localSIPDomain: 'sbc.example.com', + writeAlerts: () => {}, + realtimeDbHelpers: { + createEphemeralGateway: () => Promise.resolve(), + deleteEphemeralGateway: () => Promise.resolve() + }, + dbHelpers: { + updateVoipCarriersRegisterStatus: (sid, s) => state.statusUpdates.push(s) + } + } + }; + return { srf, state }; +} + +const ok200 = { status: 200, reason: 'OK', has: () => false, get: () => undefined, getParsedHeader: () => [] }; +const tick = () => new Promise((resolve) => setImmediate(resolve)); + +test('reregister sends a fresh REGISTER and re-arms the timer on 200 OK', async(t) => { + const { srf, state } = makeSrf(); + const rb = new Regbot(logger, REGBOT_OPTS); + // pretend a stale re-register timer is pending from before the drachtio restart + rb.timer = setTimeout(() => t.fail('stale timer should have been cleared'), 600000); + const prevEpoch = rb.epoch; + + rb.reregister(srf); + await tick(); + + t.equal(state.requests.length, 1, 'a REGISTER was sent'); + t.ok(rb.watchdog, 'response watchdog is armed while awaiting the response'); + t.ok(rb.epoch > prevEpoch, 'attempt epoch advanced'); + + state.requests[0].emit('response', ok200); + await tick(); + + t.equal(rb.status, 'registered', 'status is registered after 200 OK'); + t.equal(rb.watchdog, null, 'watchdog is cleared once the response arrives'); + t.ok(rb.timer, 'next re-register timer is armed'); + clearTimeout(rb.timer); + t.end(); +}); + +test('reregister on a retired regbot is a no-op', async(t) => { + const { srf, state } = makeSrf(); + const rb = new Regbot(logger, REGBOT_OPTS); + rb.retired = true; + + rb.reregister(srf); + await tick(); + + t.equal(state.requests.length, 0, 'no REGISTER is sent for a retired regbot'); + t.notOk(rb.timer, 'no timer is armed'); + t.end(); +}); + +test('a stale REGISTER response (superseded by a newer attempt) is ignored', async(t) => { + const { srf, state } = makeSrf(); + const rb = new Regbot(logger, REGBOT_OPTS); + + rb.register(srf); // attempt 1 + await tick(); + const req1 = state.requests[0]; + + rb.reregister(srf); // attempt 2 supersedes attempt 1 + await tick(); + const req2 = state.requests[1]; + t.ok(req2 && req2 !== req1, 'a second REGISTER was sent by reregister'); + + // a late response to attempt 1 must not be acted on + req1.emit('response', ok200); + await tick(); + t.notEqual(rb.status, 'registered', 'stale response did not mark the regbot registered'); + t.ok(rb.watchdog, 'the current attempt watchdog is still armed after the stale response'); + + // the real response to attempt 2 is processed normally + req2.emit('response', ok200); + await tick(); + t.equal(rb.status, 'registered', 'current attempt response is processed'); + t.equal(rb.watchdog, null, 'watchdog cleared by the current response'); + clearTimeout(rb.timer); + t.end(); +}); + +test('watchdog marks fail and backs off when no SIP response arrives', async(t) => { + // reload regbot with a short response timeout so the watchdog fires quickly. + // NB: this leaves the require cache holding a fresh Regbot class; harmless here since + // nothing compares class identity, but restore the cache at the end for later suites. + process.env.JAMBONES_REGBOT_RESPONSE_TIMEOUT = '40'; + clearModule('../lib/config'); + clearModule('../lib/regbot'); + const RegbotFresh = require('../lib/regbot'); + + const { srf, state } = makeSrf(); + const rb = new RegbotFresh(logger, REGBOT_OPTS); + + rb.register(srf); // no response is ever emitted + await tick(); + t.equal(state.requests.length, 1, 'first REGISTER sent'); + t.ok(rb.watchdog, 'watchdog armed while awaiting response'); + + await new Promise((resolve) => setTimeout(resolve, 90)); + t.equal(state.requests.length, 1, 'watchdog does NOT re-send inline (no REGISTER storm)'); + t.equal(rb.status, 'fail', 'watchdog marked the regbot failed'); + t.ok(rb.timer, 'a retry timer (FAILURE_RETRY_INTERVAL backoff) was scheduled'); + t.ok(state.statusUpdates.length >= 1, 'a fail status was written to the db'); + + // a late response for the timed-out attempt is ignored (epoch was bumped) + state.requests[0].emit('response', ok200); + await tick(); + t.equal(rb.status, 'fail', 'late response after the watchdog fired is ignored'); + + // stop the retry and restore the module cache for later suites + rb.retired = true; + clearTimeout(rb.timer); + clearTimeout(rb.watchdog); + clearModule('../lib/config'); + clearModule('../lib/regbot'); + delete process.env.JAMBONES_REGBOT_RESPONSE_TIMEOUT; + t.end(); +});