From b26ca68fe13b222edba9cd3184e42e61f0389ac6 Mon Sep 17 00:00:00 2001 From: = Date: Wed, 3 Sep 2025 20:22:04 +0530 Subject: [PATCH] feat: added support for redis cluster --- backend/src/@types/fastify.d.ts | 4 ++-- .../ee/services/event/event-bus-service.ts | 4 ++-- .../ee/services/event/event-sse-service.ts | 4 ++-- .../src/ee/services/event/event-sse-stream.ts | 4 ++-- backend/src/lib/config/env.ts | 20 ++++++++++++++++--- backend/src/lib/config/redis.ts | 18 ++++++++++++++++- backend/src/queue/queue-service.ts | 4 ++++ backend/src/server/app.ts | 5 +++-- 8 files changed, 49 insertions(+), 14 deletions(-) diff --git a/backend/src/@types/fastify.d.ts b/backend/src/@types/fastify.d.ts index c25d8d4d1..085c13ec8 100644 --- a/backend/src/@types/fastify.d.ts +++ b/backend/src/@types/fastify.d.ts @@ -1,6 +1,6 @@ import "fastify"; -import { Redis } from "ioredis"; +import { Cluster, Redis } from "ioredis"; import { TUsers } from "@app/db/schemas"; import { TAccessApprovalPolicyServiceFactory } from "@app/ee/services/access-approval-policy/access-approval-policy-types"; @@ -194,7 +194,7 @@ declare module "fastify" { } interface FastifyInstance { - redis: Redis; + redis: Redis | Cluster; services: { login: TAuthLoginFactory; password: TAuthPasswordFactory; diff --git a/backend/src/ee/services/event/event-bus-service.ts b/backend/src/ee/services/event/event-bus-service.ts index bb102c721..00fd14e9f 100644 --- a/backend/src/ee/services/event/event-bus-service.ts +++ b/backend/src/ee/services/event/event-bus-service.ts @@ -1,11 +1,11 @@ -import Redis from "ioredis"; +import { Cluster, Redis } from "ioredis"; import { z } from "zod"; import { logger } from "@app/lib/logger"; import { BusEventSchema, TopicName } from "./types"; -export const eventBusFactory = (redis: Redis) => { +export const eventBusFactory = (redis: Redis | Cluster) => { const publisher = redis.duplicate(); // Duplicate the publisher to create a subscriber. // This is necessary because Redis does not allow a single connection to both publish and subscribe. diff --git a/backend/src/ee/services/event/event-sse-service.ts b/backend/src/ee/services/event/event-sse-service.ts index dc52cc14c..147af8e3c 100644 --- a/backend/src/ee/services/event/event-sse-service.ts +++ b/backend/src/ee/services/event/event-sse-service.ts @@ -1,6 +1,6 @@ /* eslint-disable no-continue */ import { subject } from "@casl/ability"; -import Redis from "ioredis"; +import { Cluster, Redis } from "ioredis"; import { KeyStorePrefixes } from "@app/keystore/keystore"; import { logger } from "@app/lib/logger"; @@ -12,7 +12,7 @@ import { BusEvent, RegisteredEvent } from "./types"; const AUTH_REFRESH_INTERVAL = 60 * 1000; const HEART_BEAT_INTERVAL = 15 * 1000; -export const sseServiceFactory = (bus: TEventBusService, redis: Redis) => { +export const sseServiceFactory = (bus: TEventBusService, redis: Redis | Cluster) => { const clients = new Set(); const heartbeatInterval = setInterval(() => { diff --git a/backend/src/ee/services/event/event-sse-stream.ts b/backend/src/ee/services/event/event-sse-stream.ts index 13e18374f..9ed5e677c 100644 --- a/backend/src/ee/services/event/event-sse-stream.ts +++ b/backend/src/ee/services/event/event-sse-stream.ts @@ -3,7 +3,7 @@ import { Readable } from "node:stream"; import { MongoAbility, PureAbility } from "@casl/ability"; import { MongoQuery } from "@ucast/mongo2js"; -import Redis from "ioredis"; +import { Cluster, Redis } from "ioredis"; import { nanoid } from "nanoid"; import { ProjectType } from "@app/db/schemas"; @@ -65,7 +65,7 @@ export type EventStreamClient = { matcher: PureAbility; }; -export function createEventStreamClient(redis: Redis, options: IEventStreamClientOpts): EventStreamClient { +export function createEventStreamClient(redis: Redis | Cluster, options: IEventStreamClientOpts): EventStreamClient { const rules = options.registered.map((r) => { const secretPath = r.conditions?.secretPath; const hasConditions = r.conditions?.environmentSlug || r.conditions?.secretPath; diff --git a/backend/src/lib/config/env.ts b/backend/src/lib/config/env.ts index 586a69655..77c1ae3d9 100644 --- a/backend/src/lib/config/env.ts +++ b/backend/src/lib/config/env.ts @@ -37,6 +37,8 @@ const envSchema = z .default("false") .transform((el) => el === "true"), REDIS_URL: zpStr(z.string().optional()), + REDIS_USERNAME: zpStr(z.string().optional()), + REDIS_PASSWORD: zpStr(z.string().optional()), REDIS_SENTINEL_HOSTS: zpStr( z .string() @@ -49,6 +51,12 @@ const envSchema = z REDIS_SENTINEL_ENABLE_TLS: zodStrBool.optional().describe("Whether to use TLS/SSL for Redis Sentinel connection"), REDIS_SENTINEL_USERNAME: zpStr(z.string().optional().describe("Authentication username for Redis Sentinel")), REDIS_SENTINEL_PASSWORD: zpStr(z.string().optional().describe("Authentication password for Redis Sentinel")), + REDIS_CLUSTER_HOSTS: zpStr( + z + .string() + .optional() + .describe("Comma-separated list of Sentinel host:port pairs. Eg: 192.168.65.254:26379,192.168.65.254:26380") + ), HOST: zpStr(z.string().default("localhost")), DB_CONNECTION_URI: zpStr(z.string().describe("Postgres database connection string")).default( `postgresql://${process.env.DB_USER}:${process.env.DB_PASSWORD}@${process.env.DB_HOST}:${process.env.DB_PORT}/${process.env.DB_NAME}` @@ -335,8 +343,8 @@ const envSchema = z "Either ENCRYPTION_KEY or ROOT_ENCRYPTION_KEY must be defined." ) .refine( - (data) => Boolean(data.REDIS_URL) || Boolean(data.REDIS_SENTINEL_HOSTS), - "Either REDIS_URL or REDIS_SENTINEL_HOSTS must be defined." + (data) => Boolean(data.REDIS_URL) || Boolean(data.REDIS_SENTINEL_HOSTS) || Boolean(data.REDIS_CLUSTER_HOSTS), + "Either REDIS_URL, REDIS_SENTINEL_HOSTS or REDIS_CLUSTER_HOSTS must be defined." ) .transform((data) => ({ ...data, @@ -346,7 +354,7 @@ const envSchema = z : undefined, isCloud: Boolean(data.LICENSE_SERVER_KEY), isSmtpConfigured: Boolean(data.SMTP_HOST), - isRedisConfigured: Boolean(data.REDIS_URL || data.REDIS_SENTINEL_HOSTS), + isRedisConfigured: Boolean(data.REDIS_URL || data.REDIS_SENTINEL_HOSTS || data.REDIS_CLUSTER_HOSTS), isDevelopmentMode: data.NODE_ENV === "development", isTestMode: data.NODE_ENV === "test", isRotationDevelopmentMode: @@ -361,6 +369,12 @@ const envSchema = z const [host, port] = el.trim().split(":"); return { host: host.trim(), port: Number(port.trim()) }; }), + REDIS_CLUSTER_HOSTS: data.REDIS_CLUSTER_HOSTS?.trim() + ?.split(",") + .map((el) => { + const [host, port] = el.trim().split(":"); + return { host: host.trim(), port: Number(port.trim()) }; + }), isSecretScanningConfigured: Boolean(data.SECRET_SCANNING_GIT_APP_ID) && Boolean(data.SECRET_SCANNING_PRIVATE_KEY) && diff --git a/backend/src/lib/config/redis.ts b/backend/src/lib/config/redis.ts index 987518dd5..620f0f936 100644 --- a/backend/src/lib/config/redis.ts +++ b/backend/src/lib/config/redis.ts @@ -2,6 +2,11 @@ import { Redis } from "ioredis"; export type TRedisConfigKeys = Partial<{ REDIS_URL: string; + REDIS_USERNAME: string; + REDIS_PASSWORD: string; + + REDIS_CLUSTER_HOSTS: { host: string; port: number }[]; + REDIS_SENTINEL_HOSTS: { host: string; port: number }[]; REDIS_SENTINEL_MASTER_NAME: string; REDIS_SENTINEL_ENABLE_TLS: boolean; @@ -12,6 +17,15 @@ export type TRedisConfigKeys = Partial<{ export const buildRedisFromConfig = (cfg: TRedisConfigKeys) => { if (cfg.REDIS_URL) return new Redis(cfg.REDIS_URL, { maxRetriesPerRequest: null }); + if (cfg.REDIS_CLUSTER_HOSTS) { + return new Redis.Cluster(cfg.REDIS_CLUSTER_HOSTS, { + redisOptions: { + username: cfg.REDIS_USERNAME, + password: cfg.REDIS_PASSWORD + } + }); + } + return new Redis({ // refine at tope will catch this case sentinels: cfg.REDIS_SENTINEL_HOSTS!, @@ -19,6 +33,8 @@ export const buildRedisFromConfig = (cfg: TRedisConfigKeys) => { maxRetriesPerRequest: null, sentinelUsername: cfg.REDIS_SENTINEL_USERNAME, sentinelPassword: cfg.REDIS_SENTINEL_PASSWORD, - enableTLSForSentinelMode: cfg.REDIS_SENTINEL_ENABLE_TLS + enableTLSForSentinelMode: cfg.REDIS_SENTINEL_ENABLE_TLS, + username: cfg.REDIS_USERNAME, + password: cfg.REDIS_PASSWORD }); }; diff --git a/backend/src/queue/queue-service.ts b/backend/src/queue/queue-service.ts index f4c49d15a..9ca9a09d5 100644 --- a/backend/src/queue/queue-service.ts +++ b/backend/src/queue/queue-service.ts @@ -415,6 +415,7 @@ export const queueServiceFactory = ( redisCfg: TRedisConfigKeys, { dbConnectionUrl, dbRootCert }: { dbConnectionUrl: string; dbRootCert?: string } ): TQueueServiceFactory => { + const isClusterMode = Boolean(redisCfg?.REDIS_CLUSTER_HOSTS); const connection = buildRedisFromConfig(redisCfg); const queueContainer = {} as Record< QueueName, @@ -457,6 +458,8 @@ export const queueServiceFactory = ( } queueContainer[name] = new Queue(name as string, { + // ref: docs.bullmq.io/bull/patterns/redis-cluster + prefix: isClusterMode ? `{${name}}` : undefined, ...queueSettings, ...(crypto.isFipsModeEnabled() ? { @@ -472,6 +475,7 @@ export const queueServiceFactory = ( const appCfg = getConfig(); if (appCfg.QUEUE_WORKERS_ENABLED && isQueueEnabled(name)) { workerContainer[name] = new Worker(name, jobFn, { + prefix: isClusterMode ? `{${name}}` : undefined, ...queueSettings, ...(crypto.isFipsModeEnabled() ? { diff --git a/backend/src/server/app.ts b/backend/src/server/app.ts index 321a3656e..8cf23f703 100644 --- a/backend/src/server/app.ts +++ b/backend/src/server/app.ts @@ -12,7 +12,7 @@ import type { FastifyRateLimitOptions } from "@fastify/rate-limit"; import ratelimiter from "@fastify/rate-limit"; import { fastifyRequestContext } from "@fastify/request-context"; import fastify from "fastify"; -import { Redis } from "ioredis"; +import { Cluster, Redis } from "ioredis"; import { Knex } from "knex"; import { HsmModule } from "@app/ee/services/hsm/hsm-types"; @@ -43,7 +43,7 @@ type TMain = { queue: TQueueServiceFactory; keyStore: TKeyStoreFactory; hsmModule: HsmModule; - redis: Redis; + redis: Redis | Cluster; envConfig: TEnvConfig; superAdminDAL: TSuperAdminDALFactory; }; @@ -76,6 +76,7 @@ export const main = async ({ server.setValidatorCompiler(validatorCompiler); server.setSerializerCompiler(serializerCompiler); + // @ts-expect-error akhilmhdh: even on setting it fastify as Redis | Cluster it's throwing error server.decorate("redis", redis); server.addContentTypeParser("application/scim+json", { parseAs: "string" }, (_, body, done) => { try {