diff --git a/app.js b/app.js index ea3c25a..73a74ee 100644 --- a/app.js +++ b/app.js @@ -152,6 +152,7 @@ srf.locals.matcher = matcher; srf.connect({ host: DRACHTIO_HOST, port: DRACHTIO_PORT, secret: DRACHTIO_SECRET }); let drachtioConnected = false; +let sbcKeepAliveTimer; 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 @@ -233,19 +234,22 @@ srf.on('connect', (err, hp, version, localHostports) => { logger.info({ips: [...mapOfPublicAddresses.entries()]}, 'drachtio sip public contacts'); - mapOfPublicAddresses.forEach((addr) => { - addSbcAddress(addr.ipv4, addr.port, addr.tls_port, addr.wss_port); - // keep alive for this SBC - setTimeout(() => { - addSbcAddress(addr.ipv4, addr.port, addr.tls_port, addr.wss_port); - }, interval); - }); - - // first start up, clean sbc address - cleanSbcAddresses(); - setTimeout(() => { - cleanSbcAddresses(); - }, interval); + /* Register this SBC's public addresses and keep them alive. + addSbcAddress() refreshes the row's last_updated, and cleanSbcAddresses() (run here and by + every other sidecar in the cluster) deletes rows older than DEAD_SBC_IN_SECOND, so the + refresh MUST recur for as long as this process lives. We reap stale rows only after + refreshing our own, so the cleaner can never run ahead of this SBC's first keepalive. */ + const sbcKeepAlive = async() => { + for (const addr of mapOfPublicAddresses.values()) { + await addSbcAddress(addr.ipv4, addr.port, addr.tls_port, addr.wss_port); + } + await cleanSbcAddresses(); + }; + sbcKeepAlive(); + // drachtio re-emits 'connect' on reconnect (possibly with different contacts): replace, don't stack + if (sbcKeepAliveTimer) clearInterval(sbcKeepAliveTimer); + sbcKeepAliveTimer = setInterval(sbcKeepAlive, interval); + sbcKeepAliveTimer.unref(); /* start regbot */ require('./lib/sip-trunk-register')(logger, srf); diff --git a/lib/options.js b/lib/options.js index 3f08ddf..dd345cf 100644 --- a/lib/options.js +++ b/lib/options.js @@ -12,46 +12,48 @@ const rtpServers = new Map(); module.exports = ({srf, logger}) => { const {stats, addToSet, removeFromSet, isMemberOfSet, retrieveSet, matcher} = srf.locals; + const {client} = srf.locals.realtimeDbHelpers; const setNameFs = `${(JAMBONES_CLUSTER_ID || 'default')}:active-fs`; const setNameRtp = `${(JAMBONES_CLUSTER_ID || 'default')}:active-rtp`; const setNameFsSeriveUrl = `${(JAMBONES_CLUSTER_ID || 'default')}:fs-service-url`; + /* The active-fs / fs-service-url / active-rtp sets are shared by every SBC in the cluster, + so the decision to expire a member must be based on the last time ANY SBC heard from it, + not just this one. Each SBC records the last ping it received per member in a shared + redis hash (`:lastseen`, member -> epoch ms); the sweep below removes a member + only when that shared timestamp is stale. A member with no shared timestamp is never + expired by the sweep - e.g. an SBC that has been taken out of active-sip and is no longer + being pinged must not remove members the other SBCs are still hearing from. */ + const lastSeenHash = (setName) => `${setName}:lastseen`; + + const _expireStaleMembers = async(map, setName, now, expires) => { + const lastSeen = await client.hgetall(lastSeenHash(setName)); + for (const [key, ts] of Object.entries(lastSeen)) { + if (now - Number(ts) > expires) { + map.delete(key); + await client.hdel(lastSeenHash(setName), key); + await removeFromSet(setName, key); + const members = await retrieveSet(setName); + const countOfMembers = members.length; + logger.info({members}, `expired member ${key} from ${setName} we now have ${countOfMembers}`); + } + } + }; + /* check for expired servers every so often */ - setInterval(async() => { + const expiryTimer = setInterval(async() => { const now = Date.now(); const expires = EXPIRES_INTERVAL || 60000; - for (const [key, value] of fsServers) { - const duration = now - value; - if (duration > expires) { - fsServers.delete(key); - await removeFromSet(setNameFs, key); - const members = await retrieveSet(setNameFs); - const countOfMembers = members.length; - logger.info({members}, `expired member ${key} from ${setNameFs} we now have ${countOfMembers}`); - } - } - for (const [key, value] of fsServiceUrls) { - const duration = now - value; - if (duration > expires) { - fsServiceUrls.delete(key); - await removeFromSet(setNameFsSeriveUrl, key); - const members = await retrieveSet(setNameFsSeriveUrl); - const countOfMembers = members.length; - logger.info({members}, `expired member ${key} from ${setNameFsSeriveUrl} we now have ${countOfMembers}`); - } - } - for (const [key, value] of rtpServers) { - const duration = now - value; - if (duration > expires) { - rtpServers.delete(key); - await removeFromSet(setNameRtp, key); - const members = await retrieveSet(setNameRtp); - const countOfMembers = members.length; - logger.info({members}, `expired member ${key} from ${setNameRtp} we now have ${countOfMembers}`); - } + try { + await _expireStaleMembers(fsServers, setNameFs, now, expires); + await _expireStaleMembers(fsServiceUrls, setNameFsSeriveUrl, now, expires); + await _expireStaleMembers(rtpServers, setNameRtp, now, expires); + } catch (err) { + logger.error({err}, 'error expiring stale members'); } }, CHECK_EXPIRES_INTERVAL || 20000); + expiryTimer.unref(); /* retrieve the initial list of servers, if any, so we can watch them as well */ const _init = async() => { @@ -84,7 +86,9 @@ module.exports = ({srf, logger}) => { const _addToCache = async(map, status, setName, key) => { let countOfMembers; if (status === 'open') { - map.set(key, Date.now()); + const now = Date.now(); + map.set(key, now); + await client.hset(lastSeenHash(setName), key, now); const exists = await isMemberOfSet(setName, key); if (!exists) { await addToSet(setName, key); @@ -101,6 +105,7 @@ module.exports = ({srf, logger}) => { } else { map.delete(key); + await client.hdel(lastSeenHash(setName), key); await removeFromSet(setName, key); const members = await retrieveSet(setName); countOfMembers = members.length; diff --git a/test/index.js b/test/index.js index 019c19f..34fda45 100644 --- a/test/index.js +++ b/test/index.js @@ -6,6 +6,7 @@ require('./regbot-concurrent-rebuild-test'); require('./regbot-reconnect-test'); require('./sip-register-tests'); require('./sip-options-tests'); +require('./options-shared-expiry-test'); require('./cli-tests'); require('./docker_stop'); require('./utils'); diff --git a/test/options-shared-expiry-test.js b/test/options-shared-expiry-test.js new file mode 100644 index 0000000..70acceb --- /dev/null +++ b/test/options-shared-expiry-test.js @@ -0,0 +1,64 @@ +const test = require('tape'); +const clearModule = require('clear-module'); + +/* Unit test for the expiry sweep in lib/options.js. + The active-fs / fs-service-url / active-rtp sets are shared by every SBC in the cluster, + so a member must only be expired when the SHARED last-seen timestamp is stale, never + because this particular SBC has stopped receiving pings. */ + +const wait = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); + +test('options: expiry sweep uses shared last-seen, not local view', async(t) => { + clearModule.all(); + process.env.EXPIRES_INTERVAL = '1000'; + process.env.CHECK_EXPIRES_INTERVAL = '300'; + + const logger = require('pino')({level: 'silent'}); + const rdb = require('@jambonz/realtimedb-helpers')({}, logger); + const {client, addToSet, removeFromSet, isMemberOfSet, retrieveSet} = rdb; + const stats = {gauge: () => {}}; + const srf = {locals: {stats, addToSet, removeFromSet, isMemberOfSet, retrieveSet, realtimeDbHelpers: {client}}}; + + const setName = 'default:active-fs'; + const lastSeen = `${setName}:lastseen`; + const fresh = '10.0.0.1:5060'; // pinged another SBC recently + const stale = '10.0.0.2:5060'; // nobody has heard from it + const legacy = '10.0.0.3:5060'; // in the set with no shared timestamp + + const cleanup = async() => { + await client.del(setName, lastSeen); + }; + + try { + await cleanup(); + await client.sadd(setName, fresh, stale, legacy); + await client.hset(lastSeen, stale, Date.now() - 5000); + + /* start the options handler: this SBC has never been pinged by any of these members */ + require('../lib/options')({srf, logger}); + + /* simulate another SBC continuing to hear from `fresh` */ + const refresher = setInterval(() => client.hset(lastSeen, fresh, Date.now()), 200); + + await wait(1500); + clearInterval(refresher); + + const members = (await retrieveSet(setName)).sort(); + t.ok(members.includes(fresh), 'member still being pinged by another SBC is kept'); + t.notOk(members.includes(stale), 'member nobody has heard from within EXPIRES_INTERVAL is expired'); + t.ok(members.includes(legacy), 'member with no shared last-seen is left alone'); + t.equal(await client.hexists(lastSeen, stale), 0, 'expired member removed from shared last-seen hash'); + + await cleanup(); + client.quit(); + t.end(); + } catch (err) { + await cleanup(); + client.quit(); + t.end(err); + } finally { + delete process.env.EXPIRES_INTERVAL; + delete process.env.CHECK_EXPIRES_INTERVAL; + clearModule.all(); + } +});