diff --git a/README.md b/README.md index 4f47e402..a569627a 100644 --- a/README.md +++ b/README.md @@ -68,6 +68,19 @@ config = { } ``` +### Redis client (node-redis v6) + +Orka uses [node-redis](https://github.com/redis/node-redis) v6. The client returned by `getRedis()` +is promise-based (`await getRedis().get('key')`) — the callback API of node-redis v3 is gone. +The `config.redis` schema is unchanged: `url`, `options.tls` and the legacy retry/keepalive options +(`timesConnected`, `totalRetryTime`, `reconnectAfterMultiplier`, `socketKeepalive`, `socketInitialDelay`) +are mapped to the new driver's `socket` options and reconnect strategy. Any other key in +`config.redis.options` (e.g. `RESP`, `pingInterval`, `socket`) is passed through to `createClient`. +Orka pins `RESP: 2` by default to preserve the wire protocol reply shapes; set +`config.redis.options.RESP = 3` to opt into RESP3. +The initial connection is awaited during boot; if redis is down the app still starts +and `/health` reports unhealthy. + ```js const { orka } = require('@workablehr/orka'); diff --git a/examples/redis-example/routes.js b/examples/redis-example/routes.js index 767d65fe..d01f9c7c 100644 --- a/examples/redis-example/routes.js +++ b/examples/redis-example/routes.js @@ -1,5 +1,4 @@ const { getRedis } = require('../../build'); -const { promisify } = require('util'); const { middlewares: { health } } = require('../../build'); @@ -9,12 +8,12 @@ module.exports = { get: { health: health, '/key': async (ctx, next) => { - ctx.body = await promisify(redis.get.bind(redis))('key'); + ctx.body = await redis.get('key'); } }, put: { '/key': async (ctx, next) => { - ctx.body = await promisify(redis.set.bind(redis))('key', ctx.request.body.key); + ctx.body = await redis.set('key', ctx.request.body.key); } } }; diff --git a/package-lock.json b/package-lock.json index 1f3b6c90..805af669 100644 --- a/package-lock.json +++ b/package-lock.json @@ -27,7 +27,7 @@ "node-cron": "^2.0.3", "qs": "^6.15.3", "rabbit-queue": "^5.9.1", - "redis": "^3.1.1", + "redis": "^6.1.0", "sanitize-html": "^2.17.5", "source-map-support": "^0.5.16", "tsconfig-paths": "^3.9.0", @@ -42,7 +42,6 @@ "@types/node": "^20.7.0", "@types/pg": "^7.14.11", "@types/qs": "^6.9.7", - "@types/redis": "^2.8.14", "@types/sinon": "^21.0.0", "axios": "^1.18.1", "bullmq": "*", @@ -1377,7 +1376,7 @@ "version": "1.9.0", "resolved": "https://registry.npmjs.org/@opentelemetry/api/-/api-1.9.0.tgz", "integrity": "sha512-3giAOQvZiH5F9bMlMiv8+GSPMeqg0dbaeo58/0SlA9sxSqZhnUtxzX9/2FzyhS9sWQf5S0GJE0AKBrFqjpeYcg==", - "dev": true, + "devOptional": true, "license": "Apache-2.0", "engines": { "node": ">=8.0.0" @@ -2033,6 +2032,87 @@ "dev": true, "license": "BSD-3-Clause" }, + "node_modules/@redis/bloom": { + "version": "6.1.0", + "resolved": "https://registry.npmjs.org/@redis/bloom/-/bloom-6.1.0.tgz", + "integrity": "sha512-Rzascjd9J9bJsM45T/Z9CTg1QY/B63B6YO8QorLVMeXnbBDsKiSCVR/+GQ061hYPk8FpTzWmPY8tAv2sT+JEtQ==", + "license": "MIT", + "engines": { + "node": ">= 20.0.0" + }, + "peerDependencies": { + "@redis/client": "^6.1.0" + } + }, + "node_modules/@redis/client": { + "version": "6.1.0", + "resolved": "https://registry.npmjs.org/@redis/client/-/client-6.1.0.tgz", + "integrity": "sha512-7u1LefkezJF0HESlhO7ZFLEPfyY+NejP3SGv+Z4pGaT3oM5GVVLa0u3f4rDLUrcw+SRo8IlX9Y8JAONeDdg1Ag==", + "license": "MIT", + "dependencies": { + "cluster-key-slot": "1.1.2" + }, + "engines": { + "node": ">= 20.0.0" + }, + "peerDependencies": { + "@node-rs/xxhash": "^1.1.0", + "@opentelemetry/api": ">=1 <2" + }, + "peerDependenciesMeta": { + "@node-rs/xxhash": { + "optional": true + }, + "@opentelemetry/api": { + "optional": true + } + } + }, + "node_modules/@redis/client/node_modules/cluster-key-slot": { + "version": "1.1.2", + "resolved": "https://registry.npmjs.org/cluster-key-slot/-/cluster-key-slot-1.1.2.tgz", + "integrity": "sha512-RMr0FhtfXemyinomL4hrWcYJxmX6deFdCxpJzhDttxgO1+bcCnkk+9drydLVDmAMG7NE6aN/fl4F7ucU/90gAA==", + "license": "Apache-2.0", + "engines": { + "node": ">=0.10.0" + } + }, + "node_modules/@redis/json": { + "version": "6.1.0", + "resolved": "https://registry.npmjs.org/@redis/json/-/json-6.1.0.tgz", + "integrity": "sha512-/GFjQA6bu5pG9ClCJAI5Xx4bNXe7UTpxBBlIupBNTrn1+nY860apGnYJuaSCDV2BmEbTidpa7O2qa28oxKx+rg==", + "license": "MIT", + "engines": { + "node": ">= 20.0.0" + }, + "peerDependencies": { + "@redis/client": "^6.1.0" + } + }, + "node_modules/@redis/search": { + "version": "6.1.0", + "resolved": "https://registry.npmjs.org/@redis/search/-/search-6.1.0.tgz", + "integrity": "sha512-kS5agg+3yZbrdrt8omrew7FLCD8eOm7tarG1CROekPBRe+QGDR9aOpnHIQaYsYi6wPRTH70nQiF06AIjgURefQ==", + "license": "MIT", + "engines": { + "node": ">= 20.0.0" + }, + "peerDependencies": { + "@redis/client": "^6.1.0" + } + }, + "node_modules/@redis/time-series": { + "version": "6.1.0", + "resolved": "https://registry.npmjs.org/@redis/time-series/-/time-series-6.1.0.tgz", + "integrity": "sha512-uIDBtV8MmG/xpJsRqbGSO4iX6ryj37MLMP82lRpFvI7ykAVe5GyqgxigEbU+uZNv9kDPNMKw3dvI/S/J1BNBzA==", + "license": "MIT", + "engines": { + "node": ">= 20.0.0" + }, + "peerDependencies": { + "@redis/client": "^6.1.0" + } + }, "node_modules/@samverschueren/stream-to-observable": { "version": "0.3.1", "resolved": "https://registry.npmjs.org/@samverschueren/stream-to-observable/-/stream-to-observable-0.3.1.tgz", @@ -2357,16 +2437,6 @@ "integrity": "sha512-hKormJbkJqzQGhziax5PItDUTMAM9uE2XXQmM37dyd4hVM+5aVl7oVxMVUiVQn2oCQFN/LKCZdvSM0pFRqbSmQ==", "license": "MIT" }, - "node_modules/@types/redis": { - "version": "2.8.32", - "resolved": "https://registry.npmjs.org/@types/redis/-/redis-2.8.32.tgz", - "integrity": "sha512-7jkMKxcGq9p242exlbsVzuJb57KqHRhNl4dHoQu2Y5v9bCAbtIXXH0R3HleSQW4CTOqpHIYUW3t6tpUj4BVQ+w==", - "dev": true, - "license": "MIT", - "dependencies": { - "@types/node": "*" - } - }, "node_modules/@types/send": { "version": "1.2.1", "resolved": "https://registry.npmjs.org/@types/send/-/send-1.2.1.tgz", @@ -8945,34 +9015,26 @@ } }, "node_modules/redis": { - "version": "3.1.2", - "resolved": "https://registry.npmjs.org/redis/-/redis-3.1.2.tgz", - "integrity": "sha512-grn5KoZLr/qrRQVwoSkmzdbw6pwF+/rwODtrOr6vuBRiR/f3rjSTGupbF90Zpqm2oenix8Do6RV7pYEkGwlKkw==", + "version": "6.1.0", + "resolved": "https://registry.npmjs.org/redis/-/redis-6.1.0.tgz", + "integrity": "sha512-0kvUPM8RHP/ZMa0xYaDTcG5e8tIGW6kz6MToVT0V8iOnk6bkXp2jncGRGe2bZEk41lZwiDspUqjZCSk5ohjcKw==", "license": "MIT", "dependencies": { - "denque": "^1.5.0", - "redis-commands": "^1.7.0", - "redis-errors": "^1.2.0", - "redis-parser": "^3.0.0" + "@redis/bloom": "6.1.0", + "@redis/client": "6.1.0", + "@redis/json": "6.1.0", + "@redis/search": "6.1.0", + "@redis/time-series": "6.1.0" }, "engines": { - "node": ">=10" - }, - "funding": { - "type": "opencollective", - "url": "https://opencollective.com/node-redis" + "node": ">= 20.0.0" } }, - "node_modules/redis-commands": { - "version": "1.7.0", - "resolved": "https://registry.npmjs.org/redis-commands/-/redis-commands-1.7.0.tgz", - "integrity": "sha512-nJWqw3bTFy21hX/CPKHth6sfhZbdiHP6bTawSgQBlKOVRG7EZkfHbbHwQJnrE4vsQf0CMNE+3gJ4Fmm16vdVlQ==", - "license": "MIT" - }, "node_modules/redis-errors": { "version": "1.2.0", "resolved": "https://registry.npmjs.org/redis-errors/-/redis-errors-1.2.0.tgz", "integrity": "sha512-1qny3OExCf0UvUV/5wpYKf2YwPcOqXzkwKKSmKHiE6ZMQs5heeE/c8eXK+PNllPvmjgAbfnsbpkGZWy8cBpn9w==", + "dev": true, "license": "MIT", "engines": { "node": ">=4" @@ -8982,6 +9044,7 @@ "version": "3.0.0", "resolved": "https://registry.npmjs.org/redis-parser/-/redis-parser-3.0.0.tgz", "integrity": "sha512-DJnGAeenTdpMEH6uAJRK/uiyEIH9WVsUmoLwzudwGJUwZPp80PDBWPHXSAGNPwNvIXAbe7MSUB1zQFugFml66A==", + "dev": true, "license": "MIT", "dependencies": { "redis-errors": "^1.0.0" @@ -8990,15 +9053,6 @@ "node": ">=4" } }, - "node_modules/redis/node_modules/denque": { - "version": "1.5.1", - "resolved": "https://registry.npmjs.org/denque/-/denque-1.5.1.tgz", - "integrity": "sha512-XwE+iZ4D6ZUB7mfYRMb5wByE8L74HCn30FBN7sWnXksWc1LO1bPDl67pBR9o/kC4z/xSNAwkMYcGgqDV3BE3Hw==", - "license": "Apache-2.0", - "engines": { - "node": ">=0.10" - } - }, "node_modules/release-zalgo": { "version": "1.0.0", "resolved": "https://registry.npmjs.org/release-zalgo/-/release-zalgo-1.0.0.tgz", diff --git a/package.json b/package.json index 5c73546b..7e33c9c8 100644 --- a/package.json +++ b/package.json @@ -71,7 +71,7 @@ "node-cron": "^2.0.3", "qs": "^6.15.3", "rabbit-queue": "^5.9.1", - "redis": "^3.1.1", + "redis": "^6.1.0", "sanitize-html": "^2.17.5", "source-map-support": "^0.5.16", "tsconfig-paths": "^3.9.0", @@ -86,7 +86,6 @@ "@types/node": "^20.7.0", "@types/pg": "^7.14.11", "@types/qs": "^6.9.7", - "@types/redis": "^2.8.14", "@types/sinon": "^21.0.0", "axios": "^1.18.1", "bullmq": "*", diff --git a/src/initializers/redis.ts b/src/initializers/redis.ts index bf613372..faa7a543 100644 --- a/src/initializers/redis.ts +++ b/src/initializers/redis.ts @@ -1,9 +1,11 @@ -import { createClient as createClientType, RedisClient as RedisClientType } from 'redis'; +import type { createClient as createClientType, RedisClientType } from 'redis'; import { getLogger } from './log4js'; import { isEmpty, cloneDeep } from 'lodash'; const logger = getLogger('services.redisService'); +export type OrkaRedisClient = RedisClientType; + function getRedisUrl(config) { return config && config.url; } @@ -12,9 +14,15 @@ function getHost(url) { return url.split('@')[1] || url; } -let firstClient: RedisClientType; +function isConnectionRefused(cause) { + if (!cause) return false; + if (cause.code === 'ECONNREFUSED') return true; + return Array.isArray(cause.errors) && cause.errors.some(e => e?.code === 'ECONNREFUSED'); +} + +let firstClient: OrkaRedisClient; -export function createRedisConnection(config) { +export async function createRedisConnection(config) { const { createClient }: { createClient: typeof createClientType } = require('redis'); config = cloneDeep(config); const redisUrl = getRedisUrl(config); @@ -27,45 +35,68 @@ export function createRedisConnection(config) { if (isEmpty(config.options.tls)) delete config.options.tls; } - const options = { - timesConnected: 10, - totalRetryTime: 1000 * 60 * 60, - reconnectAfterMultiplier: 1000, - socketKeepalive: true, - socketInitialDelay: 60000, - ...config.options - }; + const { + timesConnected = 10, + totalRetryTime = 1000 * 60 * 60, + reconnectAfterMultiplier = 1000, + socketKeepalive = true, + socketInitialDelay = 60000, + tls, + socket: socketOptions, + ...clientOptions + } = config.options ?? {}; + + let timesConnectedCounter = 0; + let firstRetryAt: number; - options.retry_strategy = function(opts) { - logger.error('Retrying to connect to Redis', opts); - if (opts.error && opts.error.code === 'ECONNREFUSED') return new Error('The server refused the connection'); - if (opts.total_retry_time > options.totalRetryTime) return new Error('Retry time exhausted'); - if (opts.times_connected > options.timesConnected) { - const msg = - 'redis error retry_strategy options.times_connected exhausted.' + - 'Please verify that the redis-server "timeout" config is large enough or disabled(0).' + - 'Redis-cli:"config get timeout" '; - // This will be thrown globally and will stop the server - throw new Error(msg); + // Preserves the semantics of the v3 retry_strategy: give up on refused connections, + // on exhausted total retry time, and after too many reconnections (the health check + // then reports unhealthy so the orchestrator can restart the service). + const reconnectStrategy = (retries: number, cause: Error): number | Error => { + logger.error(`Retrying to connect to Redis (retries: ${retries})`, cause); + if (isConnectionRefused(cause)) return new Error('The server refused the connection'); + if (firstRetryAt === undefined) firstRetryAt = Date.now(); + if (Date.now() - firstRetryAt > totalRetryTime) return new Error('Retry time exhausted'); + if (timesConnectedCounter > timesConnected) { + return new Error( + 'redis error reconnectStrategy timesConnected exhausted.' + + 'Please verify that the redis-server "timeout" config is large enough or disabled(0).' + + 'Redis-cli:"config get timeout" ' + ); } - const retryInMS = Math.pow(2, opts.attempt) * options.reconnectAfterMultiplier; + const retryInMS = Math.pow(2, retries + 1) * reconnectAfterMultiplier; logger.info(`Retrying to connect to redis in ${retryInMS}ms`); return retryInMS; }; - const client = createClient(redisUrl, options); + const client = createClient({ + url: redisUrl, + RESP: 2, + ...clientOptions, + socket: { + keepAlive: socketKeepalive, + keepAliveInitialDelay: socketInitialDelay, + reconnectStrategy, + ...(tls ? { tls: true, ...tls } : {}), + ...socketOptions + } + }); if (!firstClient) firstClient = client; - client.on('connect', () => { - const socket = client.stream; - (socket as any).setKeepAlive(options.socketKeepalive, options.socketInitialDelay); + + client.on('ready', () => { + timesConnectedCounter++; + firstRetryAt = undefined; logger.info(`Redis connected ${getHost(redisUrl)}`); }); client.on('error', e => { - if (Array.isArray(e.args)) e.args[0] = getHost(redisUrl); - logger.error(e, `Redis disconnected`); + logger.error(e, `Redis disconnected ${getHost(redisUrl)}`); }); + // Awaited during orka boot so the client is ready before the server starts listening. + // A failed connection is only logged: the app still boots and /health reports unhealthy. + await client.connect().catch(e => logger.error(e, `Redis connection failed ${getHost(redisUrl)}`)); + return client; } @@ -76,5 +107,5 @@ export function getRedis() { export const isHealthy = () => { if (!firstClient) return false; - return firstClient.connected; + return firstClient.isReady; }; diff --git a/test/examples/health-example.test.ts b/test/examples/health-example.test.ts index 828b0fe5..78898f47 100644 --- a/test/examples/health-example.test.ts +++ b/test/examples/health-example.test.ts @@ -91,7 +91,7 @@ describe('Health examples', () => { }); it('/health returns not ok', async () => { - getRedis().end(true); + getRedis().destroy(); await supertest('localhost:3210').get('/health').expect(503); }); }); diff --git a/test/initializers/redis.test.ts b/test/initializers/redis.test.ts index 01bf02bc..f4757d49 100644 --- a/test/initializers/redis.test.ts +++ b/test/initializers/redis.test.ts @@ -18,14 +18,21 @@ describe('Redis connection', function() { } }; let redis; - let connectStub: sinon.SinonSpy; - let onStub: sinon.SinonSpy; + let createClientStub: sinon.SinonStub; + let onStub: sinon.SinonStub; + let connectStub: sinon.SinonStub; + + const createdOptions = () => createClientStub.args[0][0]; + const reconnectStrategy = () => createdOptions().socket.reconnectStrategy; + const emit = (event: string, ...args) => + onStub.args.filter(([e]) => e === event).forEach(([, handler]) => handler(...args)); beforeEach(async function() { onStub = sandbox.stub(); - connectStub = sandbox.stub().returns({ on: onStub }); + connectStub = sandbox.stub().resolves(); + createClientStub = sandbox.stub().returns({ on: onStub, connect: connectStub }); delete require.cache[require.resolve('../../src/initializers/redis')]; - mock('redis', { createClient: connectStub }); + mock('redis', { createClient: createClientStub }); ({ createRedisConnection: redis } = await import('../../src/initializers/redis')); }); @@ -34,40 +41,40 @@ describe('Redis connection', function() { mock.stopAll(); }); - it('should connect to redis', () => { - redis(config); - delete connectStub.args[0][1].retry_strategy; - connectStub.args.should.eql([ + it('should connect to redis', async () => { + await redis(config); + delete createdOptions().socket.reconnectStrategy; + createClientStub.args.should.eql([ [ - 'redis://localhost:6379/', { - timesConnected: 10, - totalRetryTime: 3600000, - reconnectAfterMultiplier: 1000, - socketKeepalive: true, - socketInitialDelay: 60000, - sample: 'sample' + url: 'redis://localhost:6379/', + RESP: 2, + sample: 'sample', + socket: { + keepAlive: true, + keepAliveInitialDelay: 60000 + } } ] ]); + connectStub.calledOnce.should.be.true(); }); - it('should connect to redis with tls', () => { + it('should connect to redis with tls', async () => { const newConfig = cloneDeep(config); newConfig.options.tls.key = 'key'; - redis(newConfig); - delete connectStub.args[0][1].retry_strategy; - connectStub.args.should.eql([ + await redis(newConfig); + delete createdOptions().socket.reconnectStrategy; + createClientStub.args.should.eql([ [ - 'redis://localhost:6379/', { - timesConnected: 10, - totalRetryTime: 3600000, - reconnectAfterMultiplier: 1000, - socketKeepalive: true, - socketInitialDelay: 60000, + url: 'redis://localhost:6379/', + RESP: 2, sample: 'sample', - tls: { + socket: { + keepAlive: true, + keepAliveInitialDelay: 60000, + tls: true, key: 'key' } } @@ -75,53 +82,69 @@ describe('Redis connection', function() { ]); }); - it('should connect to redis without options in config', () => { + it('should connect to redis without options in config', async () => { const newConfig = cloneDeep(config); delete newConfig.options; - redis(newConfig); - delete connectStub.args[0][1].retry_strategy; - connectStub.args.should.eql([ + await redis(newConfig); + delete createdOptions().socket.reconnectStrategy; + createClientStub.args.should.eql([ [ - 'redis://localhost:6379/', { - timesConnected: 10, - totalRetryTime: 3600000, - reconnectAfterMultiplier: 1000, - socketKeepalive: true, - socketInitialDelay: 60000 + url: 'redis://localhost:6379/', + RESP: 2, + socket: { + keepAlive: true, + keepAliveInitialDelay: 60000 + } } ] ]); }); - it('should not connect to redis', () => { - redis({}); - connectStub.args.should.eql([]); + it('should not connect to redis', async () => { + await redis({}); + createClientStub.args.should.eql([]); }); - describe('retry_strategy', () => { - it('returns server refused error', () => { - redis(config); - connectStub.args[0][1] - .retry_strategy({ error: { code: 'ECONNREFUSED' } }) - .should.eql(new Error('The server refused the connection')); + describe('reconnectStrategy', () => { + it('returns server refused error', async () => { + await redis(config); + reconnectStrategy()(0, { code: 'ECONNREFUSED' }).should.eql(new Error('The server refused the connection')); + }); + + it('returns server refused error for aggregate errors', async () => { + await redis(config); + reconnectStrategy()(0, { errors: [{ code: 'ECONNREFUSED' }] }).should.eql( + new Error('The server refused the connection') + ); + }); + + it('returns retry time exhausted error', async () => { + const clock = sandbox.useFakeTimers(); + await redis(config); + reconnectStrategy()(0, new Error('boom')).should.equal(2000); + clock.tick(1000 * 60 * 60 + 1); + reconnectStrategy()(1, new Error('boom')).should.eql(new Error('Retry time exhausted')); }); - it('returns retry time exhausted error', () => { - redis(config); - connectStub.args[0][1] - .retry_strategy({ total_retry_time: 1000 * 60 * 60 + 1 }) - .should.eql(new Error('Retry time exhausted')); + it('resets the retry time window after a successful reconnection', async () => { + const clock = sandbox.useFakeTimers(); + await redis(config); + reconnectStrategy()(0, new Error('boom')).should.equal(2000); + clock.tick(1000 * 60 * 60 + 1); + emit('ready'); + reconnectStrategy()(0, new Error('boom')).should.equal(2000); }); - it('throws error after 10 times connected - server will restart after that', () => { - redis(config); - (() => connectStub.args[0][1].retry_strategy({ times_connected: 11 })).should.throwError(); + it('returns error after 10 times connected - health check will report unhealthy after that', async () => { + await redis(config); + for (let i = 0; i < 11; i++) emit('ready'); + reconnectStrategy()(0, new Error('boom')).should.be.an.instanceOf(Error); }); - it('returns ms to retry connection', function() { - redis(config); - connectStub.args[0][1].retry_strategy({ attempt: 2 }).should.equal(4000); + it('returns ms to retry connection', async function() { + await redis(config); + reconnectStrategy()(1, new Error('boom')).should.equal(4000); }); }); });