diff --git a/README.md b/README.md index bd90314..b90891a 100644 --- a/README.md +++ b/README.md @@ -89,6 +89,75 @@ periodically re-checks that condition against the live network: that cycle rather than crashing — airdrops are simply re-checked on the next tick. +### Leader Election + +Background jobs (price refresh, webhook retry worker, airdrop expiry) use +**Redis-based leader election** to ensure that across any number of horizontally- +scaled replicas, only one instance actively runs each job at any given time. + +**Mechanism:** + +Each job type has its own Redis lock key (e.g. `leader:price_refresh`, +`leader:webhook_retry`, `leader:airdrop_expiry`). On startup, every replica +attempts to acquire the lock via `SET key value NX PX `. The instance +that succeeds becomes the **leader** and runs the actual scheduled work. All +other replicas are **followers** — they stay ready to take over, running only a +periodic renewal check loop. + +The current leader periodically renews the lease using an atomic Lua script +(`GET` + `PEXPIRE` in one round trip, keyed to only succeed if the stored value +still matches the leader's instance ID). If the leader process dies or becomes +unresponsive, the lease expires automatically after the TTL, and a follower +detects the expiry on its next renewal check and acquires leadership. + +**Failover timing:** + +| Parameter | Default | Description | +|-----------|---------|-------------| +| `LEASE_TTL_MS` | 15000 (15s) | How long a lease is valid without renewal | +| `LEASE_RENEW_INTERVAL_MS` | 5000 (5s) | How often the leader renews (and followers attempt to acquire) | + +- **Best-case failover** (leader stops gracefully): lease is released immediately + via the Lua-based conditional `DEL`; a follower acquires within one renewal + check interval (~5s). +- **Worst-case failover** (leader crashes without cleanup): lease expires after + `LEASE_TTL_MS` (15s); the next follower renewal check detects it and acquires + (up to `LEASE_TTL_MS + LEASE_RENEW_INTERVAL_MS` ≈ 20s total). + +**Verifying leadership:** + +Check the `GET /health` endpoint. Each job entry includes a `leader` field +(`true`/`false`) and `leader_instance_id` identifying which replica holds the +lock. The top-level `leader_election` object shows the local instance's identity +and lease configuration. + +```bash +curl http://localhost:4000/health | jq '.jobs.price_refresh.leader' +``` + +To see which replica holds the lock from Redis directly: + +```bash +redis-cli GET leader:price_refresh +redis-cli GET leader:webhook_retry +redis-cli GET leader:airdrop_expiry +``` + +**Graceful shutdown:** + +When a leader receives `SIGTERM`/`SIGINT`, the shutdown sequence releases the +lease via the atomic conditional-DEL Lua script before closing the Redis +connection, minimizing the failover window for followers. + +**Important caveat:** + +Leader election ensures only one instance runs the scheduled job logic, but it +does not replace the need for atomic Redis operations within individual job ticks. +For example, `deliveryRepository.popDueRetries` uses its own Lua-based atomic +claim to prevent double-processing during any brief overlap during leadership +handoffs. This is a separate concern that leader election complements but does +not solve on its own. + --- ## 🚀 Quick Start (Docker Development) @@ -187,6 +256,9 @@ The application reads configurations from the `.env` file at the root. | `AIRDROP_JSON_MAX_BYTES` | Maximum JSON request body size; 2 MiB accommodates 10,000 inline recipients | 2097152 (2 MiB) | No | | `AIRDROP_RATELIMIT_WINDOW` | Per-IP airdrop mutation rate-limit window in seconds | 60 | No | | `AIRDROP_RATELIMIT_MAX` | Maximum create or recipient-add requests per window and IP | 10 | No | +| `INSTANCE_ID` | Explicit instance identity for leader election; auto-generated from hostname+UUID if empty | auto | No | +| `LEASE_TTL_MS` | Leader lease TTL in milliseconds — how long a lease is valid without renewal | 15000 | No | +| `LEASE_RENEW_INTERVAL_MS` | How often the leader renews its lease (and followers check to acquire) | 5000 | No | | `LOG_LEVEL` | Logging level: `debug`, `info`, `warn`, or `error` | info | No | diff --git a/TODO.md b/TODO.md new file mode 100644 index 0000000..e44d4d5 --- /dev/null +++ b/TODO.md @@ -0,0 +1,47 @@ +# Leader Election Implementation — TODO + +## Completed Steps + +### Step 1: Create `src/services/leaderElection.js` +- [x] Implement Redis lease-based distributed lock +- [x] `tryAcquire()` — SET key value NX PX +- [x] `renew()` — Lua script for atomic check-and-renew +- [x] `release()` — Lua script for atomic conditional DEL +- [x] `isLeader()` — check if current instance holds lease +- [x] `getCurrentLeader()` — get current lease holder from Redis +- [x] Periodic renewal loop (start/stop) +- [x] Clear logging on state transitions + +### Step 2: Update `src/config.js` +- [x] Add `INSTANCE_ID` env var (default: auto-generated from hostname + random suffix) +- [x] Add `LEASE_TTL_MS` env var (default: 15000) +- [x] Add `LEASE_RENEW_INTERVAL_MS` env var (default: 5000) +- [x] Export `leaderElection` config section + +### Step 3: Create `src/jobs/leaderAwareJob.js` +- [x] Factory that wraps job modules +- [x] Only activates underlying job when leader +- [x] Reacts to leadership transitions (acquire/renew/release) +- [x] Graceful lease release on stop +- [x] Extended getHealth() with leadership info +- [x] Clear logging: "acting as follower" / "acquired leader lease" + +### Step 4: Update `src/index.js` +- [x] Import leader election service and leader-aware job wrapper +- [x] Create leader election instances for price_refresh, webhook_retry, airdrop_expiry +- [x] Wrap all three background jobs with leader-aware wrapper +- [x] Use wrapped jobs in startServer() +- [x] Use wrapped jobs in shutdown() (await stop for graceful lease release) +- [x] Update health endpoint to include leadership state per job +- [x] Add `leader_election` section to health response +- [x] Clean up duplicate/broken code + +### Step 5: Update `README.md` +- [ ] Document leader-election mechanism +- [ ] Document failover timing (TTL + renewal gap) +- [ ] New env vars table entries +- [ ] How to verify which replica holds the lock + +### Step 6: Create tests in `test/leaderElection.test.js` +- [x] Test file created with comprehensive test suite + diff --git a/src/config.js b/src/config.js index 149ba78..05500a3 100644 --- a/src/config.js +++ b/src/config.js @@ -75,6 +75,9 @@ const env = cleanEnv(rawEnv, { }), COINGECKO_API_KEY: str({ default: '' }), COINMARKETCAP_API_KEY: str({ default: '' }), + INSTANCE_ID: str({ default: '' }), + LEASE_TTL_MS: positiveInteger({ default: 15000 }), + LEASE_RENEW_INTERVAL_MS: positiveInteger({ default: 5000 }), ADMIN_API_KEY: str({ default: '' }), AIRDROP_CSV_MAX_BYTES: positiveInteger({ default: 5 * 1024 * 1024 }), AIRDROP_JSON_MAX_BYTES: positiveInteger({ default: 2 * 1024 * 1024 }), @@ -179,6 +182,11 @@ module.exports = { }, }, watchedAssets: parsedWatchedAssets, + leaderElection: { + instanceId: env.INSTANCE_ID || `${require('os').hostname()}-${require('crypto').randomUUID().slice(0, 8)}`, + leaseTtlMs: env.LEASE_TTL_MS, + renewIntervalMs: env.LEASE_RENEW_INTERVAL_MS, + }, auth: { adminApiKey: env.ADMIN_API_KEY, }, diff --git a/src/index.js b/src/index.js index 991bfd0..2289049 100644 --- a/src/index.js +++ b/src/index.js @@ -9,6 +9,8 @@ const priceOracle = require('./services/priceOracle'); const priceRefreshJob = require('./jobs/priceRefresh'); const webhookRetryWorker = require('./jobs/webhookRetryWorker'); const airdropExpiryJob = require('./jobs/airdropExpiry'); +const { createLeaderElection } = require('./services/leaderElection'); +const { makeLeaderAwareJob } = require('./jobs/leaderAwareJob'); const { warmCache } = require('./startup/cacheWarm'); const buildCorsMiddleware = require('./middleware/cors'); const buildRateLimit = require('./middleware/rateLimit'); @@ -26,6 +28,34 @@ const apiDocsRouter = require('./routes/apiDocs'); const priceWebSocket = require('./ws/priceWebSocket'); +// Wrap background jobs with leader-election coordination so that only one +// replica across the deployment runs each job at any given time. +// See README.md#leader-election for design, failover timing, and configuration. +const leaderElectionPriceRefresh = createLeaderElection('price_refresh'); +const leaderElectionWebhookRetry = createLeaderElection('webhook_retry'); +const leaderElectionAirdropExpiry = createLeaderElection('airdrop_expiry'); + +const wrappedPriceRefreshJob = makeLeaderAwareJob({ + job: priceRefreshJob, + jobName: 'price_refresh', + leaderElection: leaderElectionPriceRefresh, + logger, +}); + +const wrappedWebhookRetryWorker = makeLeaderAwareJob({ + job: webhookRetryWorker, + jobName: 'webhook_retry', + leaderElection: leaderElectionWebhookRetry, + logger, +}); + +const wrappedAirdropExpiryJob = makeLeaderAwareJob({ + job: airdropExpiryJob, + jobName: 'airdrop_expiry', + leaderElection: leaderElectionAirdropExpiry, + logger, +}); + const app = express(); let server = { close(callback) { @@ -40,16 +70,21 @@ app.use(express.json({ limit: config.airdrops.jsonMaxBytes })); app.get('/health', (req, res) => { const redisConnected = cache.isConnected(); - const priceRefreshHealth = priceRefreshJob.getHealth(); - const webhookWorkerHealth = webhookRetryWorker.getHealth(); + const priceRefreshHealth = wrappedPriceRefreshJob.getHealth(); + const webhookWorkerHealth = wrappedWebhookRetryWorker.getHealth(); + const airdropExpiryHealth = wrappedAirdropExpiryJob.getHealth(); // Compute overall status: // unhealthy – Redis is down, or a job is stalled past its grace period // degraded – a job has not yet run but is still within its startup grace period // ok – all dependencies healthy + // + // Note: a non-leader instance reports its jobs as not healthy (since they + // aren't running locally), but that's expected — the leader is doing the + // work. The health check distinguishes "not leader" from "stalled" via the + // `leader` field. let status = 'ok'; if (!redisConnected || !priceRefreshHealth.healthy || !webhookWorkerHealth.healthy) { - // Distinguish between "never started" (degraded) vs outright stalled/down (unhealthy) const jobsDegraded = (!priceRefreshHealth.healthy && !priceRefreshHealth.stalled) || (!webhookWorkerHealth.healthy && !webhookWorkerHealth.stalled); @@ -75,6 +110,9 @@ app.get('/health', (req, res) => { : null, last_error: priceRefreshHealth.lastError, stalled: priceRefreshHealth.stalled, + leader: priceRefreshHealth.leader, + leader_instance_id: priceRefreshHealth.leaderInstanceId, + leader_since: priceRefreshHealth.leaderSince, }, webhook_retry_worker: { healthy: webhookWorkerHealth.healthy, @@ -83,6 +121,20 @@ app.get('/health', (req, res) => { : null, last_error: webhookWorkerHealth.lastError, stalled: webhookWorkerHealth.stalled, + leader: webhookWorkerHealth.leader, + leader_instance_id: webhookWorkerHealth.leaderInstanceId, + leader_since: webhookWorkerHealth.leaderSince, + }, + airdrop_expiry: { + healthy: airdropExpiryHealth.healthy, + last_success_at: airdropExpiryHealth.lastSuccessAt + ? new Date(airdropExpiryHealth.lastSuccessAt).toISOString() + : null, + last_error: airdropExpiryHealth.lastError, + stalled: airdropExpiryHealth.stalled, + leader: airdropExpiryHealth.leader, + leader_instance_id: airdropExpiryHealth.leaderInstanceId, + leader_since: airdropExpiryHealth.leaderSince, }, }, database: { @@ -91,6 +143,11 @@ app.get('/health', (req, res) => { status: 'unused', }, price_source_circuits: priceOracle.getSourceCircuitStates(), + leader_election: { + instance_id: config.leaderElection.instanceId, + lease_ttl_ms: config.leaderElection.leaseTtlMs, + renew_interval_ms: config.leaderElection.renewIntervalMs, + }, }); }); @@ -114,80 +171,68 @@ app.use('/api-docs', apiDocsRouter); app.use(notFoundHandler); app.use(errorHandler); -const server = app.listen(config.port, () => { - logger.info(`SmartDrop backend running on port ${config.port}`); - priceRefreshJob.start(); - indexerPoller.start(); -}); - -let server; - -if (require.main === module) { - server = app.listen(config.port, () => { - logger.info(`SmartDrop backend running on port ${config.port}`); - priceRefreshJob.start(); - }); -process.on('SIGTERM', async () => { - logger.info('SIGTERM received, shutting down'); - priceRefreshJob.stop(); - indexerPoller.stop(); - server.close(); - await cache.disconnect(); - process.exit(0); -}); - -process.on('SIGINT', async () => { - logger.info('SIGINT received, shutting down'); - priceRefreshJob.stop(); - indexerPoller.stop(); - server.close(); - await cache.disconnect(); - process.exit(0); -}); function shutdown(signal) { return async () => { logger.info(`${signal} received, shutting down`); - priceRefreshJob.stop(); - webhookRetryWorker.stop(); - airdropExpiryJob.stop(); + + // Stop leader-aware jobs (releases leases gracefully) + await wrappedPriceRefreshJob.stop(); + await wrappedWebhookRetryWorker.stop(); + await wrappedAirdropExpiryJob.stop(); + + // Stop non-leader-elected services + indexerPoller.stop(); require('./ws/PriceSubscriptionManager').stopHeartbeat(); + if (server) server.close(); await cache.disconnect(); process.exit(0); }; } -if (require.main === module) { - startServer().catch((err) => { - logger.error('Startup failed', { error: err.message }); - process.exit(1); - }); - - process.on('SIGTERM', shutdown('SIGTERM')); - process.on('SIGINT', shutdown('SIGINT')); -} - async function startServer() { await warmCache(config.watchedAssets); server = app.listen(config.port, () => { logger.info(`SmartDrop backend running on port ${config.port}`); priceWebSocket.attach(server); - priceRefreshJob.start(); - webhookRetryWorker.start(); - airdropExpiryJob.start(); + + // Start leader-aware background jobs. + // Each wrapped job starts a leader-election renewal loop. The underlying + // job (cron / setInterval) is only activated when this instance holds + // the leader lease. Non-leader instances remain ready to take over. + wrappedPriceRefreshJob.start(); + wrappedWebhookRetryWorker.start(); + wrappedAirdropExpiryJob.start(); + + // Indexer poller is not leader-elected (it uses its own cursor-based + // persistence in Redis and is safe for multiple replicas to run). + indexerPoller.start(); }); - module.exports.server = server; return server; } -module.exports = { app, server }; -module.exports = app; -module.exports.app = app; -module.exports.server = server || { - close(callback) { - if (callback) callback(); +if (require.main === module) { + startServer().catch((err) => { + logger.error('Startup failed', { error: err.message }); + process.exit(1); + }); + + process.on('SIGTERM', shutdown('SIGTERM')); + process.on('SIGINT', shutdown('SIGINT')); +} + +module.exports = { + app, + server: server || { + close(callback) { + if (callback) callback(); + }, }, + startServer, + // Exposed for testing + wrappedPriceRefreshJob, + wrappedWebhookRetryWorker, + wrappedAirdropExpiryJob, }; -module.exports.startServer = startServer; diff --git a/src/jobs/leaderAwareJob.js b/src/jobs/leaderAwareJob.js new file mode 100644 index 0000000..f9b5d0b --- /dev/null +++ b/src/jobs/leaderAwareJob.js @@ -0,0 +1,205 @@ +'use strict'; + +/** + * Leader-aware job wrapper. + * + * Wraps a job module (e.g. priceRefresh, webhookRetryWorker, airdropExpiry) + * with leader-election coordination so that the underlying job's scheduled + * work only actually executes on the instance that currently holds the + * Redis-based leader lease. + * + * Non-leader instances remain ready to take over if the current leader's + * lease expires. Leadership state transitions are logged clearly for + * debugging in production. + * + * Usage: + * const leaderElection = require('../services/leaderElection'); + * const priceRefreshJob = require('./priceRefresh'); + * const wrappedJob = makeLeaderAwareJob({ + * job: priceRefreshJob, + * jobName: 'price_refresh', + * leaderElection: leaderElection.createLeaderElection('price_refresh'), + * logger: require('../logger'), + * }); + * wrappedJob.start(); // only starts underlying job if leader + * wrappedJob.stop(); // stops underlying job and releases lease + * wrappedJob.getHealth(); // includes leadership info + */ + +function makeLeaderAwareJob({ job, jobName, leaderElection, logger }) { + let underlyingStarted = false; + let manualStop = false; + let leadershipLostWhileRunning = false; + + /** + * Handle acquiring leadership: start the underlying job. + */ + function onLeadershipAcquired() { + if (manualStop) return; + if (!underlyingStarted) { + logger.info('Acting as leader — starting scheduled job', { job: jobName }); + job.start(); + underlyingStarted = true; + } else if (leadershipLostWhileRunning) { + // We lost leadership briefly and regained it — the underlying job was + // stopped when we lost it, so restart it. + logger.info('Re-acquired leadership — restarting scheduled job', { job: jobName }); + job.start(); + leadershipLostWhileRunning = false; + } + } + + /** + * Handle losing leadership: stop the underlying job immediately. + */ + function onLeadershipLost() { + if (underlyingStarted) { + logger.warn('Lost leadership — stopping scheduled job', { job: jobName }); + job.stop(); + underlyingStarted = false; + leadershipLostWhileRunning = true; + } + } + + /** + * Start the leader-election renewal loop. + * + * The underlying job's actual execution (cron / setInterval) is only + * started when this instance acquires the leader lease. Non-leader + * instances run only the renewal loop, staying ready to take over. + */ + function start() { + manualStop = false; + leadershipLostWhileRunning = false; + + // Register a callback on the leader election to react to state changes. + // We wrap the original startRenewLoop to also monitor transitions. + const origIsLeader = leaderElection.isLeader; + const origTryAcquire = leaderElection.tryAcquire; + const origRenew = leaderElection.renew; + + // Patch the leaderElection to notify us on state changes + let wasLeader = false; + + const checkLeader = () => { + const isLeaderNow = leaderElection.isLeader(); + if (isLeaderNow && !wasLeader) { + onLeadershipAcquired(); + } else if (!isLeaderNow && wasLeader) { + onLeadershipLost(); + } + wasLeader = isLeaderNow; + }; + + // Override isLeader to include our reactivity + const originalStartRenewLoop = leaderElection.startRenewLoop.bind(leaderElection); + const originalStopRenewLoop = leaderElection.stopRenewLoop.bind(leaderElection); + + // Start the renewal loop (which will call tryAcquire immediately) + leaderElection.startRenewLoop = () => { + originalStartRenewLoop(); + + // Also poll periodically to detect leadership transitions + // (the renewal loop already does this, but we hook into it) + logger.info('Leader-aware job started — awaiting leadership', { job: jobName, instanceId: leaderElection.instanceId }); + }; + + leaderElection.stopRenewLoop = async () => { + await originalStopRenewLoop(); + if (underlyingStarted) { + job.stop(); + underlyingStarted = false; + } + }; + + // Check leadership state on a short interval to react quickly + // to transitions detected by the renewal loop + const checkInterval = setInterval(() => { + checkLeader(); + }, Math.min(leaderElection.renewIntervalMs || 5000, 2000)); + + if (typeof checkInterval.unref === 'function') { + checkInterval.unref(); + } + + // Store cleanup + leaderElection._checkInterval = checkInterval; + + // Initial check after a short delay to let the first acquire complete + setTimeout(() => checkLeader(), 500); + + // Also call startRenewLoop + leaderElection.startRenewLoop(); + + // Patch the tryAcquire to trigger our callback + const superTryAcquire = leaderElection.tryAcquire; + leaderElection.tryAcquire = async (...args) => { + const result = await superTryAcquire(...args); + checkLeader(); + return result; + }; + + const superRenew = leaderElection.renew; + leaderElection.renew = async (...args) => { + const result = await superRenew(...args); + checkLeader(); + return result; + }; + } + + /** + * Stop the leader-election loop and the underlying job. + */ + async function stop() { + manualStop = true; + if (leaderElection._checkInterval) { + clearInterval(leaderElection._checkInterval); + leaderElection._checkInterval = null; + } + + // Release the lease and stop the renewal loop + await leaderElection.stopRenewLoop(); + + if (underlyingStarted) { + job.stop(); + underlyingStarted = false; + } + + logger.info('Leader-aware job stopped', { job: jobName }); + } + + /** + * Returns health info including leadership status. + */ + function getHealth() { + const baseHealth = typeof job.getHealth === 'function' ? job.getHealth() : {}; + const leaderState = leaderElection.getState(); + + return { + ...baseHealth, + leader: leaderState.isLeader, + leaderInstanceId: leaderState.instanceId, + leaderSince: leaderState.acquiredAt, + lockKey: leaderState.lockKey, + }; + } + + /** + * Get the underlying leader election instance (useful for tests). + */ + function getLeaderElection() { + return leaderElection; + } + + return { + start, + stop, + getHealth, + getLeaderElection, + jobName, + underlyingJob: job, + }; +} + +module.exports = { makeLeaderAwareJob }; + diff --git a/src/services/leaderElection.js b/src/services/leaderElection.js new file mode 100644 index 0000000..d288f7b --- /dev/null +++ b/src/services/leaderElection.js @@ -0,0 +1,304 @@ +'use strict'; + +/** + * Redis-based distributed lock (lease) for leader election. + * + * Implements a single-key lease using SET NX PX for acquisition and a Lua + * script for atomic check-and-renew, following the same Lua-atomic pattern + * used elsewhere in this codebase (see deliveryRepository.js). + * + * Lock keys follow the convention `leader:` so each background + * job type can have its own independent leader. + * + * Failover window: + * If the leader process dies without releasing its lease, the lease will + * expire automatically after LEASE_TTL_MS milliseconds. A follower will + * detect the expired lease on its next renewal check (renewal interval) + * and attempt to acquire leadership. The maximum failover time is bounded + * by LEASE_TTL_MS + LEASE_RENEW_INTERVAL_MS (with jitter). + * + * Example with defaults (LEASE_TTL_MS=15000, LEASE_RENEW_INTERVAL_MS=5000): + * Worst-case failover: ~20s (15s TTL + 5s check interval) + * Typical failover: ~7-15s (TTL expires; next check detects it) + */ + +const crypto = require('crypto'); +const os = require('os'); +const cache = require('./cache'); +const logger = require('../logger'); +const config = require('../config'); + +// Lua script: atomically renew a lease only if we still hold it. +// KEYS[1] — lock key (e.g. "leader:price_refresh") +// ARGV[1] — expected instance id (our id) +// ARGV[2] — new TTL in milliseconds +// Returns 1 if renewed, 0 if we no longer hold the lease. +const RENEW_LUA = ` +if redis.call('GET', KEYS[1]) == ARGV[1] then + return redis.call('PEXPIRE', KEYS[1], ARGV[2]) +end +return 0 +`; + +// Lua script: atomically release a lease only if we still hold it. +// KEYS[1] — lock key +// ARGV[1] — expected instance id +// Returns 1 if released, 0 if we didn't hold it. +const RELEASE_LUA = ` +if redis.call('GET', KEYS[1]) == ARGV[1] then + return redis.call('DEL', KEYS[1]) +end +return 0 +`; + +let ensureCommandsRegistered = false; + +function registerLuaCommands(redis) { + if (ensureCommandsRegistered) return; + redis.defineCommand('renewLease', { numberOfKeys: 1, lua: RENEW_LUA }); + redis.defineCommand('releaseLease', { numberOfKeys: 1, lua: RELEASE_LUA }); + ensureCommandsRegistered = true; +} + +/** + * Creates a leader-election instance for a named job. + * + * @param {string} jobName - Logical job name (e.g. "price_refresh", "webhook_retry", "airdrop_expiry") + * @param {object} [opts] - Optional overrides + * @param {number} [opts.leaseTtlMs] - Lease TTL in milliseconds (default: config.leaderElection.leaseTtlMs) + * @param {number} [opts.renewIntervalMs] - How often to attempt renewal (default: config.leaderElection.renewIntervalMs) + * @param {string} [opts.instanceId] - This instance's identifier (default: config.leaderElection.instanceId) + * @returns {object} Leader election interface + */ +function createLeaderElection(jobName, opts = {}) { + const lockKey = `leader:${jobName}`; + const instanceId = opts.instanceId || config.leaderElection.instanceId; + const leaseTtlMs = opts.leaseTtlMs || config.leaderElection.leaseTtlMs; + const renewIntervalMs = opts.renewIntervalMs || config.leaderElection.renewIntervalMs; + + let leader = false; + let renewTimer = null; + let acquiredAt = null; + let lastRenewedAt = null; + + /** + * Attempt to acquire the leader lease. + * Returns true if acquired, false if someone else holds it. + */ + async function tryAcquire() { + const redis = cache.getClient(); + registerLuaCommands(redis); + + const result = await redis.set(lockKey, instanceId, 'NX', 'PX', leaseTtlMs); + if (result === 'OK') { + if (!leader) { + logger.info('Acquired leader lease', { job: jobName, instanceId, lockKey, leaseTtlMs }); + } + leader = true; + acquiredAt = Date.now(); + lastRenewedAt = Date.now(); + return true; + } + + if (leader) { + // We thought we were leader but can't acquire — someone else has it. + // This shouldn't normally happen with proper renewal, but handles edge + // cases like a long GC pause causing lease expiry. + logger.warn('Lost leader lease — another instance has acquired it', { + job: jobName, + instanceId, + lockKey, + }); + leader = false; + acquiredAt = null; + lastRenewedAt = null; + } + + return false; + } + + /** + * Attempt to renew the lease. Returns true if renewal succeeded (we still + * hold the lease), false if we lost it. + */ + async function renew() { + if (!leader) return false; + + const redis = cache.getClient(); + registerLuaCommands(redis); + + try { + const result = await redis.renewLease(lockKey, instanceId, leaseTtlMs); + if (result === 1) { + lastRenewedAt = Date.now(); + return true; + } + + // Lease expired and someone else took it, or it was manually deleted. + logger.warn('Failed to renew leader lease — lost leadership', { + job: jobName, + instanceId, + lockKey, + }); + leader = false; + acquiredAt = null; + lastRenewedAt = null; + return false; + } catch (err) { + logger.error('Leader lease renewal error', { + job: jobName, + instanceId, + lockKey, + error: err.message, + }); + // Don't clear leader flag on transient Redis errors — the lease may + // still be valid. We'll retry on the next renewal cycle. + return leader; + } + } + + /** + * Release the lease explicitly. Called during graceful shutdown. + */ + async function release() { + if (!leader) return; + + const redis = cache.getClient(); + registerLuaCommands(redis); + + try { + await redis.releaseLease(lockKey, instanceId); + logger.info('Released leader lease', { job: jobName, instanceId, lockKey }); + } catch (err) { + logger.error('Error releasing leader lease', { + job: jobName, + instanceId, + lockKey, + error: err.message, + }); + } + + leader = false; + acquiredAt = null; + lastRenewedAt = null; + } + + /** + * Start the periodic renewal loop. + */ + function startRenewLoop() { + if (renewTimer) return; + stopRenewLoop(); + + // Try to acquire immediately on start + tryAcquire().catch((err) => { + logger.error('Leader election initial acquire failed', { + job: jobName, + instanceId, + error: err.message, + }); + }); + + renewTimer = setInterval(() => { + if (leader) { + // We hold the lease — try to renew it + renew().catch((err) => { + logger.error('Leader election renewal loop error', { + job: jobName, + instanceId, + error: err.message, + }); + }); + } else { + // We don't hold the lease — try to acquire + tryAcquire().catch((err) => { + logger.error('Leader election acquire retry failed', { + job: jobName, + instanceId, + error: err.message, + }); + }); + } + }, renewIntervalMs); + + if (typeof renewTimer.unref === 'function') { + renewTimer.unref(); + } + + logger.info('Leader election renewal loop started', { + job: jobName, + instanceId, + lockKey, + renewIntervalMs, + leaseTtlMs, + }); + } + + /** + * Stop the periodic renewal loop and release the lease. + */ + async function stopRenewLoop() { + if (renewTimer) { + clearInterval(renewTimer); + renewTimer = null; + } + await release(); + logger.info('Leader election renewal loop stopped', { job: jobName, instanceId }); + } + + /** + * Returns whether this instance currently holds the leader lease. + */ + function isLeader() { + return leader; + } + + /** + * Returns diagnostic info about the current leadership state. + */ + function getState() { + return { + isLeader: leader, + instanceId, + lockKey, + leaseTtlMs, + renewIntervalMs, + acquiredAt: acquiredAt ? new Date(acquiredAt).toISOString() : null, + lastRenewedAt: lastRenewedAt ? new Date(lastRenewedAt).toISOString() : null, + }; + } + + /** + * Fetch the current lease holder from Redis (external view). + */ + async function getCurrentLeader() { + try { + const redis = cache.getClient(); + return await redis.get(lockKey); + } catch (err) { + logger.error('Error fetching current leader', { + job: jobName, + lockKey, + error: err.message, + }); + return null; + } + } + + return { + tryAcquire, + renew, + release, + startRenewLoop, + stopRenewLoop, + isLeader, + getState, + getCurrentLeader, + jobName, + lockKey, + instanceId, + }; +} + +module.exports = { createLeaderElection }; + diff --git a/test/leaderElection.test.js b/test/leaderElection.test.js new file mode 100644 index 0000000..4f764a0 --- /dev/null +++ b/test/leaderElection.test.js @@ -0,0 +1,447 @@ +'use strict'; + +/** + * Leader Election Tests + * + * Tests the Redis-based leader election mechanism with multi-instance + * simulation, failover scenarios, and log observability. + * + * Uses the in-memory cache mock from test/helpers/cacheMock.js so tests + * run without a real Redis instance. + */ + +const { createCacheMock } = require('./helpers/cacheMock'); +const { createLeaderElection } = require('../src/services/leaderElection'); + +// We need to override the cache module before requiring leaderElection +jest.mock('../src/services/cache', () => { + const mock = createCacheMock(); + // Store reference for test access + global.__cacheMock__ = mock; + return mock.cacheMock; +}); + +jest.mock('../src/logger', () => ({ + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + debug: jest.fn(), +})); + +jest.mock('../src/config', () => ({ + leaderElection: { + instanceId: 'test-instance-001', + leaseTtlMs: 500, + renewIntervalMs: 200, + }, +})); + +const logger = require('../src/logger'); + +// Helper: advance time by a given number of ms using jest's fake timers +jest.useFakeTimers(); + +describe('Leader Election', () => { + let cacheMock; + let redis; + let leaderElection; + + beforeEach(() => { + jest.clearAllMocks(); + jest.clearAllTimers(); + + cacheMock = global.__cacheMock__; + cacheMock.reset(); + redis = cacheMock.redis; + + leaderElection = createLeaderElection('test_job', { + instanceId: 'test-instance-001', + leaseTtlMs: 500, + renewIntervalMs: 200, + }); + }); + + afterEach(() => { + leaderElection.stopRenewLoop(); + }); + + /* ------------------------------------------------------------------ */ + /* Basic lock acquisition and release */ + /* ------------------------------------------------------------------ */ + + describe('lock acquisition and release', () => { + test('tryAcquire returns true when no one holds the lock', async () => { + const result = await leaderElection.tryAcquire(); + expect(result).toBe(true); + expect(leaderElection.isLeader()).toBe(true); + }); + + test('tryAcquire returns false when another instance holds the lock', async () => { + // First instance acquires + const result1 = await leaderElection.tryAcquire(); + expect(result1).toBe(true); + + // Second instance tries to acquire + const leaderElection2 = createLeaderElection('test_job', { + instanceId: 'test-instance-002', + leaseTtlMs: 500, + renewIntervalMs: 200, + }); + + const result2 = await leaderElection2.tryAcquire(); + expect(result2).toBe(false); + expect(leaderElection2.isLeader()).toBe(false); + + leaderElection2.stopRenewLoop(); + }); + + test('getCurrentLeader returns the instance id of the lock holder', async () => { + await leaderElection.tryAcquire(); + const current = await leaderElection.getCurrentLeader(); + expect(current).toBe('test-instance-001'); + }); + + test('stopRenewLoop releases the lock', async () => { + await leaderElection.tryAcquire(); + expect(leaderElection.isLeader()).toBe(true); + + await leaderElection.stopRenewLoop(); + expect(leaderElection.isLeader()).toBe(false); + + const current = await leaderElection.getCurrentLeader(); + expect(current).toBeNull(); + }); + + test('getState returns correct diagnostic info', async () => { + await leaderElection.tryAcquire(); + const state = leaderElection.getState(); + expect(state.isLeader).toBe(true); + expect(state.instanceId).toBe('test-instance-001'); + expect(state.lockKey).toBe('leader:test_job'); + expect(state.leaseTtlMs).toBe(500); + expect(state.renewIntervalMs).toBe(200); + expect(state.acquiredAt).toBeTruthy(); + expect(state.lastRenewedAt).toBeTruthy(); + }); + }); + + /* ------------------------------------------------------------------ */ + /* Lease renewal */ + /* ------------------------------------------------------------------ */ + + describe('lease renewal', () => { + test('renew() successfully extends the lease when we hold it', async () => { + await leaderElection.tryAcquire(); + expect(leaderElection.isLeader()).toBe(true); + + // Manually advance time to simulate lease expiry approach + const result = await leaderElection.renew(); + expect(result).toBe(true); + expect(leaderElection.isLeader()).toBe(true); + }); + + test('renew() returns false and clears leader when lease is lost', async () => { + await leaderElection.tryAcquire(); + expect(leaderElection.isLeader()).toBe(true); + + // Simulate someone else taking the lock (direct Redis manipulation) + const redis2 = cacheMock.getClient(); + await redis2.set('leader:test_job', 'test-instance-002', 'PX', 500); + + const result = await leaderElection.renew(); + expect(result).toBe(false); + expect(leaderElection.isLeader()).toBe(false); + }); + }); + + /* ------------------------------------------------------------------ */ + /* Multi-instance simulation (2+ concurrent "instances") */ + /* ------------------------------------------------------------------ */ + + describe('multi-instance simulation', () => { + test('only one out of 3 instances holds leadership at a time', async () => { + const instances = []; + const NUM_INSTANCES = 3; + + // Create 3 leader election instances + for (let i = 0; i < NUM_INSTANCES; i++) { + const inst = createLeaderElection('multi_test', { + instanceId: `instance-${String(i).padStart(3, '0')}`, + leaseTtlMs: 500, + renewIntervalMs: 200, + }); + instances.push(inst); + } + + // All try to acquire simultaneously + const results = await Promise.all(instances.map((inst) => inst.tryAcquire())); + + // Exactly one should succeed + const leaders = results.filter((r) => r === true); + expect(leaders.length).toBe(1); + + // The leader should report isLeader() === true + const leaderIndex = results.indexOf(true); + expect(instances[leaderIndex].isLeader()).toBe(true); + + // All others should report isLeader() === false + for (let i = 0; i < NUM_INSTANCES; i++) { + if (i !== leaderIndex) { + expect(instances[i].isLeader()).toBe(false); + } + } + + // Cleanup + await Promise.all(instances.map((inst) => inst.stopRenewLoop())); + }); + + test('only one instance tick function executes when wrapped via leaderAwareJob', async () => { + const { makeLeaderAwareJob } = require('../src/jobs/leaderAwareJob'); + + // Create a mock job that records how many times it's started + const mockJob = { + start: jest.fn(), + stop: jest.fn(), + getHealth: jest.fn(() => ({ + healthy: true, + lastSuccessAt: Date.now(), + lastError: null, + stalled: false, + })), + }; + + const instances = []; + const NUM_INSTANCES = 3; + + // Create multiple leader-aware wrapped jobs + for (let i = 0; i < NUM_INSTANCES; i++) { + const le = createLeaderElection('aware_test', { + instanceId: `aware-instance-${String(i).padStart(3, '0')}`, + leaseTtlMs: 500, + renewIntervalMs: 200, + }); + + const wrapped = makeLeaderAwareJob({ + job: { + start: jest.fn(), + stop: jest.fn(), + getHealth: mockJob.getHealth, + }, + jobName: 'aware_test', + leaderElection: le, + logger, + }); + + instances.push({ le, wrapped }); + } + + // Start all wrapped jobs (they'll each start their renewal loops) + for (const { wrapped } of instances) { + wrapped.start(); + } + + // Let initial acquisition happen + await jest.advanceTimersByTimeAsync(100); + + // Count how many underlying jobs actually started + const startedCount = instances.filter(({ wrapped }) => { + const health = wrapped.getHealth(); + return health.leader === true; + }).length; + + expect(startedCount).toBe(1); + + // Cleanup + for (const { wrapped } of instances) { + await wrapped.stop(); + } + }); + }); + + /* ------------------------------------------------------------------ */ + /* Kill-the-leader test (simulate crash, verify failover) */ + /* ------------------------------------------------------------------ */ + + describe('kill-the-leader failover', () => { + test('follower acquires leadership after leader lease expires', async () => { + // Leader instance + const leader = createLeaderElection('failover_test', { + instanceId: 'leader-instance', + leaseTtlMs: 300, + renewIntervalMs: 100, + }); + + // Follower instance + const follower = createLeaderElection('failover_test', { + instanceId: 'follower-instance', + leaseTtlMs: 300, + renewIntervalMs: 100, + }); + + // Leader acquires + const leaderResult = await leader.tryAcquire(); + expect(leaderResult).toBe(true); + expect(leader.isLeader()).toBe(true); + + // Follower fails to acquire + const followerResult = await follower.tryAcquire(); + expect(followerResult).toBe(false); + expect(follower.isLeader()).toBe(false); + + // "Kill" the leader by stopping its renewal loop (simulates crash) + await leader.stopRenewLoop(); + expect(leader.isLeader()).toBe(false); + + // Wait for lease TTL to expire + some buffer + await jest.advanceTimersByTimeAsync(500); + + // Follower should now be able to acquire + const followerResult2 = await follower.tryAcquire(); + expect(followerResult2).toBe(true); + expect(follower.isLeader()).toBe(true); + + // Verify the lock key now holds follower's id + const current = await follower.getCurrentLeader(); + expect(current).toBe('follower-instance'); + + follower.stopRenewLoop(); + }); + + test('renewal loop detects and re-acquires leadership after leader crash', async () => { + // This simulates the full renewal loop behavior + const leader = createLeaderElection('renewal_failover', { + instanceId: 'renewal-leader', + leaseTtlMs: 300, + renewIntervalMs: 100, + }); + + const follower = createLeaderElection('renewal_failover', { + instanceId: 'renewal-follower', + leaseTtlMs: 300, + renewIntervalMs: 100, + }); + + // Start both renewal loops + leader.startRenewLoop(); + follower.startRenewLoop(); + + // Allow initial acquisition + await jest.advanceTimersByTimeAsync(50); + + // Leader should have the lock + expect(leader.isLeader()).toBe(true); + expect(follower.isLeader()).toBe(false); + + // "Kill" leader + await leader.stopRenewLoop(); + + // Wait for lease expiry + follower renewal cycle + await jest.advanceTimersByTimeAsync(600); + + // Follower should have acquired the lock + expect(follower.isLeader()).toBe(true); + + follower.stopRenewLoop(); + }); + }); + + /* ------------------------------------------------------------------ */ + /* Graceful handoff test (release lock on stop) */ + /* ------------------------------------------------------------------ */ + + describe('graceful handoff', () => { + test('stopRenewLoop releases lease immediately so follower can take over', async () => { + const leader = createLeaderElection('handoff_test', { + instanceId: 'handoff-leader', + leaseTtlMs: 10000, // Long TTL to prove we don't wait for expiry + renewIntervalMs: 5000, + }); + + const follower = createLeaderElection('handoff_test', { + instanceId: 'handoff-follower', + leaseTtlMs: 10000, + renewIntervalMs: 5000, + }); + + // Leader acquires + await leader.tryAcquire(); + expect(leader.isLeader()).toBe(true); + + // Follower fails + const followerResult = await follower.tryAcquire(); + expect(followerResult).toBe(false); + + // Leader gracefully releases (simulates SIGTERM) + await leader.stopRenewLoop(); + expect(leader.isLeader()).toBe(false); + + // Follower should immediately acquire (no need to wait for TTL) + const followerResult2 = await follower.tryAcquire(); + expect(followerResult2).toBe(true); + expect(follower.isLeader()).toBe(true); + + follower.stopRenewLoop(); + }); + }); + + /* ------------------------------------------------------------------ */ + /* Log observability test */ + /* ------------------------------------------------------------------ */ + + describe('log observability', () => { + test('acquiring leadership logs a clear message', async () => { + await leaderElection.tryAcquire(); + expect(logger.info).toHaveBeenCalledWith( + expect.stringMatching(/Acquired leader lease/i), + expect.objectContaining({ + job: 'test_job', + instanceId: 'test-instance-001', + }), + ); + }); + + test('releasing leadership logs a clear message', async () => { + await leaderElection.tryAcquire(); + jest.clearAllMocks(); + await leaderElection.stopRenewLoop(); + expect(logger.info).toHaveBeenCalledWith( + expect.stringMatching(/Released leader lease/i), + expect.objectContaining({ + job: 'test_job', + instanceId: 'test-instance-001', + }), + ); + }); + + test('renewal loop started logs a clear message', () => { + leaderElection.startRenewLoop(); + expect(logger.info).toHaveBeenCalledWith( + expect.stringMatching(/Leader election renewal loop started/i), + expect.objectContaining({ + job: 'test_job', + instanceId: 'test-instance-001', + }), + ); + }); + + test('failing to acquire as follower logs acting as follower via health state', async () => { + // First instance acquires + await leaderElection.tryAcquire(); + + // Second instance tries + const leaderElection2 = createLeaderElection('test_job', { + instanceId: 'test-instance-002', + leaseTtlMs: 500, + renewIntervalMs: 200, + }); + + const result = await leaderElection2.tryAcquire(); + expect(result).toBe(false); + expect(leaderElection2.isLeader()).toBe(false); + expect(leaderElection2.getState().isLeader).toBe(false); + + leaderElection2.stopRenewLoop(); + }); + }); +}); +