mirror of
https://github.com/jambonz/sbc-sip-sidecar.git
synced 2026-10-04 02:04:20 +00:00
fix: shared-roster expiry and sbc_addresses keepalive (two-SBC cutover bugs) (#153)
* fix: expire FS/RTP roster members on shared last-seen, not one SBC's local view The active-fs, fs-service-url and active-rtp redis sets are shared by every SBC in the cluster, but the expiry sweep in lib/options.js removed members based solely on when THIS sidecar last received an OPTIONS ping from them. An SBC that was taken out of active-sip (so FS/RTP servers stopped pinging it) therefore deleted every feature server and rtpengine from the shared rosters 60s later, rejecting all inbound calls until the other SBC re-added them on its next ping cycle. Each SBC now records the last ping it received per member in a shared redis hash (<setName>:lastseen, member -> epoch ms) and the sweep removes a member only when that shared timestamp is older than EXPIRES_INTERVAL. Members with no shared timestamp are never expired by the sweep, so a parked or mixed-version SBC cannot remove members the others are still hearing from. Adds a unit test that reproduces the failure against the old code. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * fix: make sbc_addresses keepalive recurring instead of a single setTimeout addSbcAddress() refreshes the row's last_updated and cleanSbcAddresses() deletes rows older than DEAD_SBC_IN_SECOND (default 3600s), but app.js only re-called addSbcAddress once, 15 minutes after connecting. The row then went stale, and the next sidecar to start anywhere in the cluster deleted the healthy SBC's row, so new sip realms were provisioned with one SBC IP instead of two. Run one recurring timer (SBC_PUBLIC_ADDRESS_KEEP_ALIVE_IN_MILISECOND, default 15 min) that refreshes this SBC's rows and only then reaps stale ones, so the cleaner never runs ahead of this process's own keepalive. The timer is unref'd and replaced (not stacked) on drachtio reconnect. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Fable 5.1
parent
a4b9487623
commit
b51362ac04
@@ -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);
|
||||
|
||||
+36
-31
@@ -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 (`<setName>: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;
|
||||
|
||||
@@ -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');
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
});
|
||||
Reference in New Issue
Block a user