feat: added support for redis cluster

This commit is contained in:
=
2025-09-03 20:22:04 +05:30
parent 19a5f52d20
commit b26ca68fe1
8 changed files with 49 additions and 14 deletions
+2 -2
View File
@@ -1,6 +1,6 @@
import "fastify"; import "fastify";
import { Redis } from "ioredis"; import { Cluster, Redis } from "ioredis";
import { TUsers } from "@app/db/schemas"; import { TUsers } from "@app/db/schemas";
import { TAccessApprovalPolicyServiceFactory } from "@app/ee/services/access-approval-policy/access-approval-policy-types"; import { TAccessApprovalPolicyServiceFactory } from "@app/ee/services/access-approval-policy/access-approval-policy-types";
@@ -194,7 +194,7 @@ declare module "fastify" {
} }
interface FastifyInstance { interface FastifyInstance {
redis: Redis; redis: Redis | Cluster;
services: { services: {
login: TAuthLoginFactory; login: TAuthLoginFactory;
password: TAuthPasswordFactory; password: TAuthPasswordFactory;
@@ -1,11 +1,11 @@
import Redis from "ioredis"; import { Cluster, Redis } from "ioredis";
import { z } from "zod"; import { z } from "zod";
import { logger } from "@app/lib/logger"; import { logger } from "@app/lib/logger";
import { BusEventSchema, TopicName } from "./types"; import { BusEventSchema, TopicName } from "./types";
export const eventBusFactory = (redis: Redis) => { export const eventBusFactory = (redis: Redis | Cluster) => {
const publisher = redis.duplicate(); const publisher = redis.duplicate();
// Duplicate the publisher to create a subscriber. // Duplicate the publisher to create a subscriber.
// This is necessary because Redis does not allow a single connection to both publish and subscribe. // This is necessary because Redis does not allow a single connection to both publish and subscribe.
@@ -1,6 +1,6 @@
/* eslint-disable no-continue */ /* eslint-disable no-continue */
import { subject } from "@casl/ability"; import { subject } from "@casl/ability";
import Redis from "ioredis"; import { Cluster, Redis } from "ioredis";
import { KeyStorePrefixes } from "@app/keystore/keystore"; import { KeyStorePrefixes } from "@app/keystore/keystore";
import { logger } from "@app/lib/logger"; import { logger } from "@app/lib/logger";
@@ -12,7 +12,7 @@ import { BusEvent, RegisteredEvent } from "./types";
const AUTH_REFRESH_INTERVAL = 60 * 1000; const AUTH_REFRESH_INTERVAL = 60 * 1000;
const HEART_BEAT_INTERVAL = 15 * 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<EventStreamClient>(); const clients = new Set<EventStreamClient>();
const heartbeatInterval = setInterval(() => { const heartbeatInterval = setInterval(() => {
@@ -3,7 +3,7 @@ import { Readable } from "node:stream";
import { MongoAbility, PureAbility } from "@casl/ability"; import { MongoAbility, PureAbility } from "@casl/ability";
import { MongoQuery } from "@ucast/mongo2js"; import { MongoQuery } from "@ucast/mongo2js";
import Redis from "ioredis"; import { Cluster, Redis } from "ioredis";
import { nanoid } from "nanoid"; import { nanoid } from "nanoid";
import { ProjectType } from "@app/db/schemas"; import { ProjectType } from "@app/db/schemas";
@@ -65,7 +65,7 @@ export type EventStreamClient = {
matcher: PureAbility; 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 rules = options.registered.map((r) => {
const secretPath = r.conditions?.secretPath; const secretPath = r.conditions?.secretPath;
const hasConditions = r.conditions?.environmentSlug || r.conditions?.secretPath; const hasConditions = r.conditions?.environmentSlug || r.conditions?.secretPath;
+17 -3
View File
@@ -37,6 +37,8 @@ const envSchema = z
.default("false") .default("false")
.transform((el) => el === "true"), .transform((el) => el === "true"),
REDIS_URL: zpStr(z.string().optional()), REDIS_URL: zpStr(z.string().optional()),
REDIS_USERNAME: zpStr(z.string().optional()),
REDIS_PASSWORD: zpStr(z.string().optional()),
REDIS_SENTINEL_HOSTS: zpStr( REDIS_SENTINEL_HOSTS: zpStr(
z z
.string() .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_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_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_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")), HOST: zpStr(z.string().default("localhost")),
DB_CONNECTION_URI: zpStr(z.string().describe("Postgres database connection string")).default( 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}` `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." "Either ENCRYPTION_KEY or ROOT_ENCRYPTION_KEY must be defined."
) )
.refine( .refine(
(data) => Boolean(data.REDIS_URL) || Boolean(data.REDIS_SENTINEL_HOSTS), (data) => Boolean(data.REDIS_URL) || Boolean(data.REDIS_SENTINEL_HOSTS) || Boolean(data.REDIS_CLUSTER_HOSTS),
"Either REDIS_URL or REDIS_SENTINEL_HOSTS must be defined." "Either REDIS_URL, REDIS_SENTINEL_HOSTS or REDIS_CLUSTER_HOSTS must be defined."
) )
.transform((data) => ({ .transform((data) => ({
...data, ...data,
@@ -346,7 +354,7 @@ const envSchema = z
: undefined, : undefined,
isCloud: Boolean(data.LICENSE_SERVER_KEY), isCloud: Boolean(data.LICENSE_SERVER_KEY),
isSmtpConfigured: Boolean(data.SMTP_HOST), 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", isDevelopmentMode: data.NODE_ENV === "development",
isTestMode: data.NODE_ENV === "test", isTestMode: data.NODE_ENV === "test",
isRotationDevelopmentMode: isRotationDevelopmentMode:
@@ -361,6 +369,12 @@ const envSchema = z
const [host, port] = el.trim().split(":"); const [host, port] = el.trim().split(":");
return { host: host.trim(), port: Number(port.trim()) }; 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: isSecretScanningConfigured:
Boolean(data.SECRET_SCANNING_GIT_APP_ID) && Boolean(data.SECRET_SCANNING_GIT_APP_ID) &&
Boolean(data.SECRET_SCANNING_PRIVATE_KEY) && Boolean(data.SECRET_SCANNING_PRIVATE_KEY) &&
+17 -1
View File
@@ -2,6 +2,11 @@ import { Redis } from "ioredis";
export type TRedisConfigKeys = Partial<{ export type TRedisConfigKeys = Partial<{
REDIS_URL: string; 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_HOSTS: { host: string; port: number }[];
REDIS_SENTINEL_MASTER_NAME: string; REDIS_SENTINEL_MASTER_NAME: string;
REDIS_SENTINEL_ENABLE_TLS: boolean; REDIS_SENTINEL_ENABLE_TLS: boolean;
@@ -12,6 +17,15 @@ export type TRedisConfigKeys = Partial<{
export const buildRedisFromConfig = (cfg: TRedisConfigKeys) => { export const buildRedisFromConfig = (cfg: TRedisConfigKeys) => {
if (cfg.REDIS_URL) return new Redis(cfg.REDIS_URL, { maxRetriesPerRequest: null }); 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({ return new Redis({
// refine at tope will catch this case // refine at tope will catch this case
sentinels: cfg.REDIS_SENTINEL_HOSTS!, sentinels: cfg.REDIS_SENTINEL_HOSTS!,
@@ -19,6 +33,8 @@ export const buildRedisFromConfig = (cfg: TRedisConfigKeys) => {
maxRetriesPerRequest: null, maxRetriesPerRequest: null,
sentinelUsername: cfg.REDIS_SENTINEL_USERNAME, sentinelUsername: cfg.REDIS_SENTINEL_USERNAME,
sentinelPassword: cfg.REDIS_SENTINEL_PASSWORD, 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
}); });
}; };
+4
View File
@@ -415,6 +415,7 @@ export const queueServiceFactory = (
redisCfg: TRedisConfigKeys, redisCfg: TRedisConfigKeys,
{ dbConnectionUrl, dbRootCert }: { dbConnectionUrl: string; dbRootCert?: string } { dbConnectionUrl, dbRootCert }: { dbConnectionUrl: string; dbRootCert?: string }
): TQueueServiceFactory => { ): TQueueServiceFactory => {
const isClusterMode = Boolean(redisCfg?.REDIS_CLUSTER_HOSTS);
const connection = buildRedisFromConfig(redisCfg); const connection = buildRedisFromConfig(redisCfg);
const queueContainer = {} as Record< const queueContainer = {} as Record<
QueueName, QueueName,
@@ -457,6 +458,8 @@ export const queueServiceFactory = (
} }
queueContainer[name] = new Queue(name as string, { queueContainer[name] = new Queue(name as string, {
// ref: docs.bullmq.io/bull/patterns/redis-cluster
prefix: isClusterMode ? `{${name}}` : undefined,
...queueSettings, ...queueSettings,
...(crypto.isFipsModeEnabled() ...(crypto.isFipsModeEnabled()
? { ? {
@@ -472,6 +475,7 @@ export const queueServiceFactory = (
const appCfg = getConfig(); const appCfg = getConfig();
if (appCfg.QUEUE_WORKERS_ENABLED && isQueueEnabled(name)) { if (appCfg.QUEUE_WORKERS_ENABLED && isQueueEnabled(name)) {
workerContainer[name] = new Worker(name, jobFn, { workerContainer[name] = new Worker(name, jobFn, {
prefix: isClusterMode ? `{${name}}` : undefined,
...queueSettings, ...queueSettings,
...(crypto.isFipsModeEnabled() ...(crypto.isFipsModeEnabled()
? { ? {
+3 -2
View File
@@ -12,7 +12,7 @@ import type { FastifyRateLimitOptions } from "@fastify/rate-limit";
import ratelimiter from "@fastify/rate-limit"; import ratelimiter from "@fastify/rate-limit";
import { fastifyRequestContext } from "@fastify/request-context"; import { fastifyRequestContext } from "@fastify/request-context";
import fastify from "fastify"; import fastify from "fastify";
import { Redis } from "ioredis"; import { Cluster, Redis } from "ioredis";
import { Knex } from "knex"; import { Knex } from "knex";
import { HsmModule } from "@app/ee/services/hsm/hsm-types"; import { HsmModule } from "@app/ee/services/hsm/hsm-types";
@@ -43,7 +43,7 @@ type TMain = {
queue: TQueueServiceFactory; queue: TQueueServiceFactory;
keyStore: TKeyStoreFactory; keyStore: TKeyStoreFactory;
hsmModule: HsmModule; hsmModule: HsmModule;
redis: Redis; redis: Redis | Cluster;
envConfig: TEnvConfig; envConfig: TEnvConfig;
superAdminDAL: TSuperAdminDALFactory; superAdminDAL: TSuperAdminDALFactory;
}; };
@@ -76,6 +76,7 @@ export const main = async ({
server.setValidatorCompiler(validatorCompiler); server.setValidatorCompiler(validatorCompiler);
server.setSerializerCompiler(serializerCompiler); server.setSerializerCompiler(serializerCompiler);
// @ts-expect-error akhilmhdh: even on setting it fastify as Redis | Cluster it's throwing error
server.decorate("redis", redis); server.decorate("redis", redis);
server.addContentTypeParser("application/scim+json", { parseAs: "string" }, (_, body, done) => { server.addContentTypeParser("application/scim+json", { parseAs: "string" }, (_, body, done) => {
try { try {