You've already forked node-redis
mirror of
https://github.com/redis/node-redis.git
synced 2025-08-06 02:15:48 +03:00
V5 bringing RESP3, Sentinel and TypeMapping to node-redis
RESP3 Support - Some commands responses in RESP3 aren't stable yet and therefore return an "untyped" ReplyUnion. Sentinel TypeMapping Correctly types Multi commands Note: some API changes to be further documented in v4-to-v5.md
This commit is contained in:
@@ -1,260 +1,256 @@
|
||||
import { createConnection } from 'net';
|
||||
import { once } from 'events';
|
||||
import RedisClient from '@redis/client/dist/lib/client';
|
||||
import { promiseTimeout } from '@redis/client/dist/lib/utils';
|
||||
import { ClusterSlotsReply } from '@redis/client/dist/lib/commands/CLUSTER_SLOTS';
|
||||
import * as path from 'path';
|
||||
import { promisify } from 'util';
|
||||
import { exec } from 'child_process';
|
||||
import { createConnection } from 'node:net';
|
||||
import { once } from 'node:events';
|
||||
import { createClient } from '@redis/client/index';
|
||||
import { setTimeout } from 'node:timers/promises';
|
||||
// import { ClusterSlotsReply } from '@redis/client/dist/lib/commands/CLUSTER_SLOTS';
|
||||
import { promisify } from 'node:util';
|
||||
import { exec } from 'node:child_process';
|
||||
const execAsync = promisify(exec);
|
||||
|
||||
interface ErrorWithCode extends Error {
|
||||
code: string;
|
||||
code: string;
|
||||
}
|
||||
|
||||
async function isPortAvailable(port: number): Promise<boolean> {
|
||||
try {
|
||||
const socket = createConnection({ port });
|
||||
await once(socket, 'connect');
|
||||
socket.end();
|
||||
} catch (err) {
|
||||
if (err instanceof Error && (err as ErrorWithCode).code === 'ECONNREFUSED') {
|
||||
return true;
|
||||
}
|
||||
try {
|
||||
const socket = createConnection({ port });
|
||||
await once(socket, 'connect');
|
||||
socket.end();
|
||||
} catch (err) {
|
||||
if (err instanceof Error && (err as ErrorWithCode).code === 'ECONNREFUSED') {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
return false;
|
||||
return false;
|
||||
}
|
||||
|
||||
const portIterator = (async function*(): AsyncIterableIterator<number> {
|
||||
for (let i = 6379; i < 65535; i++) {
|
||||
if (await isPortAvailable(i)) {
|
||||
yield i;
|
||||
}
|
||||
const portIterator = (async function* (): AsyncIterableIterator<number> {
|
||||
for (let i = 6379; i < 65535; i++) {
|
||||
if (await isPortAvailable(i)) {
|
||||
yield i;
|
||||
}
|
||||
}
|
||||
|
||||
throw new Error('All ports are in use');
|
||||
throw new Error('All ports are in use');
|
||||
})();
|
||||
|
||||
export interface RedisServerDockerConfig {
|
||||
image: string;
|
||||
version: string;
|
||||
image: string;
|
||||
version: string;
|
||||
}
|
||||
|
||||
export interface RedisServerDocker {
|
||||
port: number;
|
||||
dockerId: string;
|
||||
port: number;
|
||||
dockerId: string;
|
||||
}
|
||||
|
||||
// ".." cause it'll be in `./dist`
|
||||
const DOCKER_FODLER_PATH = path.join(__dirname, '../docker');
|
||||
|
||||
async function spawnRedisServerDocker({ image, version }: RedisServerDockerConfig, serverArguments: Array<string>): Promise<RedisServerDocker> {
|
||||
const port = (await portIterator.next()).value,
|
||||
{ stdout, stderr } = await execAsync(
|
||||
'docker run -d --network host $(' +
|
||||
`docker build ${DOCKER_FODLER_PATH} -q ` +
|
||||
`--build-arg IMAGE=${image}:${version} ` +
|
||||
`--build-arg REDIS_ARGUMENTS="--save '' --port ${port.toString()} ${serverArguments.join(' ')}"` +
|
||||
')'
|
||||
);
|
||||
const port = (await portIterator.next()).value,
|
||||
{ stdout, stderr } = await execAsync(
|
||||
`docker run -e REDIS_ARGS="--port ${port.toString()} ${serverArguments.join(' ')}" -d --network host ${image}:${version}`
|
||||
);
|
||||
|
||||
if (!stdout) {
|
||||
throw new Error(`docker run error - ${stderr}`);
|
||||
}
|
||||
if (!stdout) {
|
||||
throw new Error(`docker run error - ${stderr}`);
|
||||
}
|
||||
|
||||
while (await isPortAvailable(port)) {
|
||||
await promiseTimeout(50);
|
||||
}
|
||||
while (await isPortAvailable(port)) {
|
||||
await setTimeout(50);
|
||||
}
|
||||
|
||||
return {
|
||||
port,
|
||||
dockerId: stdout.trim()
|
||||
};
|
||||
return {
|
||||
port,
|
||||
dockerId: stdout.trim()
|
||||
};
|
||||
}
|
||||
|
||||
const RUNNING_SERVERS = new Map<Array<string>, ReturnType<typeof spawnRedisServerDocker>>();
|
||||
|
||||
export function spawnRedisServer(dockerConfig: RedisServerDockerConfig, serverArguments: Array<string>): Promise<RedisServerDocker> {
|
||||
const runningServer = RUNNING_SERVERS.get(serverArguments);
|
||||
if (runningServer) {
|
||||
return runningServer;
|
||||
}
|
||||
const runningServer = RUNNING_SERVERS.get(serverArguments);
|
||||
if (runningServer) {
|
||||
return runningServer;
|
||||
}
|
||||
|
||||
const dockerPromise = spawnRedisServerDocker(dockerConfig, serverArguments);
|
||||
RUNNING_SERVERS.set(serverArguments, dockerPromise);
|
||||
return dockerPromise;
|
||||
const dockerPromise = spawnRedisServerDocker(dockerConfig, serverArguments);
|
||||
RUNNING_SERVERS.set(serverArguments, dockerPromise);
|
||||
return dockerPromise;
|
||||
}
|
||||
|
||||
async function dockerRemove(dockerId: string): Promise<void> {
|
||||
const { stderr } = await execAsync(`docker rm -f ${dockerId}`);
|
||||
if (stderr) {
|
||||
throw new Error(`docker rm error - ${stderr}`);
|
||||
}
|
||||
const { stderr } = await execAsync(`docker rm -f ${dockerId}`);
|
||||
if (stderr) {
|
||||
throw new Error(`docker rm error - ${stderr}`);
|
||||
}
|
||||
}
|
||||
|
||||
after(() => {
|
||||
return Promise.all(
|
||||
[...RUNNING_SERVERS.values()].map(async dockerPromise =>
|
||||
await dockerRemove((await dockerPromise).dockerId)
|
||||
)
|
||||
);
|
||||
return Promise.all(
|
||||
[...RUNNING_SERVERS.values()].map(async dockerPromise =>
|
||||
await dockerRemove((await dockerPromise).dockerId)
|
||||
)
|
||||
);
|
||||
});
|
||||
|
||||
export interface RedisClusterDockersConfig extends RedisServerDockerConfig {
|
||||
numberOfMasters?: number;
|
||||
numberOfReplicas?: number;
|
||||
numberOfMasters?: number;
|
||||
numberOfReplicas?: number;
|
||||
}
|
||||
|
||||
async function spawnRedisClusterNodeDockers(
|
||||
dockersConfig: RedisClusterDockersConfig,
|
||||
serverArguments: Array<string>,
|
||||
fromSlot: number,
|
||||
toSlot: number
|
||||
dockersConfig: RedisClusterDockersConfig,
|
||||
serverArguments: Array<string>,
|
||||
fromSlot: number,
|
||||
toSlot: number
|
||||
) {
|
||||
const range: Array<number> = [];
|
||||
for (let i = fromSlot; i < toSlot; i++) {
|
||||
range.push(i);
|
||||
}
|
||||
const range: Array<number> = [];
|
||||
for (let i = fromSlot; i < toSlot; i++) {
|
||||
range.push(i);
|
||||
}
|
||||
|
||||
const master = await spawnRedisClusterNodeDocker(
|
||||
dockersConfig,
|
||||
serverArguments
|
||||
);
|
||||
const master = await spawnRedisClusterNodeDocker(
|
||||
dockersConfig,
|
||||
serverArguments
|
||||
);
|
||||
|
||||
await master.client.clusterAddSlots(range);
|
||||
await master.client.clusterAddSlots(range);
|
||||
|
||||
if (!dockersConfig.numberOfReplicas) return [master];
|
||||
|
||||
const replicasPromises: Array<ReturnType<typeof spawnRedisClusterNodeDocker>> = [];
|
||||
for (let i = 0; i < (dockersConfig.numberOfReplicas ?? 0); i++) {
|
||||
replicasPromises.push(
|
||||
spawnRedisClusterNodeDocker(dockersConfig, [
|
||||
...serverArguments,
|
||||
'--cluster-enabled',
|
||||
'yes',
|
||||
'--cluster-node-timeout',
|
||||
'5000'
|
||||
]).then(async replica => {
|
||||
await replica.client.clusterMeet('127.0.0.1', master.docker.port);
|
||||
if (!dockersConfig.numberOfReplicas) return [master];
|
||||
|
||||
while ((await replica.client.clusterSlots()).length === 0) {
|
||||
await promiseTimeout(50);
|
||||
}
|
||||
const replicasPromises: Array<ReturnType<typeof spawnRedisClusterNodeDocker>> = [];
|
||||
for (let i = 0; i < (dockersConfig.numberOfReplicas ?? 0); i++) {
|
||||
replicasPromises.push(
|
||||
spawnRedisClusterNodeDocker(dockersConfig, [
|
||||
...serverArguments,
|
||||
'--cluster-enabled',
|
||||
'yes',
|
||||
'--cluster-node-timeout',
|
||||
'5000'
|
||||
]).then(async replica => {
|
||||
await replica.client.clusterMeet('127.0.0.1', master.docker.port);
|
||||
|
||||
await replica.client.clusterReplicate(
|
||||
await master.client.clusterMyId()
|
||||
);
|
||||
while ((await replica.client.clusterSlots()).length === 0) {
|
||||
await setTimeout(50);
|
||||
}
|
||||
|
||||
return replica;
|
||||
})
|
||||
await replica.client.clusterReplicate(
|
||||
await master.client.clusterMyId()
|
||||
);
|
||||
}
|
||||
|
||||
return [
|
||||
master,
|
||||
...await Promise.all(replicasPromises)
|
||||
];
|
||||
return replica;
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
return [
|
||||
master,
|
||||
...await Promise.all(replicasPromises)
|
||||
];
|
||||
}
|
||||
|
||||
async function spawnRedisClusterNodeDocker(
|
||||
dockersConfig: RedisClusterDockersConfig,
|
||||
serverArguments: Array<string>
|
||||
dockersConfig: RedisClusterDockersConfig,
|
||||
serverArguments: Array<string>
|
||||
) {
|
||||
const docker = await spawnRedisServerDocker(dockersConfig, [
|
||||
...serverArguments,
|
||||
'--cluster-enabled',
|
||||
'yes',
|
||||
'--cluster-node-timeout',
|
||||
'5000'
|
||||
]),
|
||||
client = RedisClient.create({
|
||||
socket: {
|
||||
port: docker.port
|
||||
}
|
||||
});
|
||||
const docker = await spawnRedisServerDocker(dockersConfig, [
|
||||
...serverArguments,
|
||||
'--cluster-enabled',
|
||||
'yes',
|
||||
'--cluster-node-timeout',
|
||||
'5000'
|
||||
]),
|
||||
client = createClient({
|
||||
socket: {
|
||||
port: docker.port
|
||||
}
|
||||
});
|
||||
|
||||
await client.connect();
|
||||
await client.connect();
|
||||
|
||||
return {
|
||||
docker,
|
||||
client
|
||||
};
|
||||
return {
|
||||
docker,
|
||||
client
|
||||
};
|
||||
}
|
||||
|
||||
const SLOTS = 16384;
|
||||
|
||||
async function spawnRedisClusterDockers(
|
||||
dockersConfig: RedisClusterDockersConfig,
|
||||
serverArguments: Array<string>
|
||||
dockersConfig: RedisClusterDockersConfig,
|
||||
serverArguments: Array<string>
|
||||
): Promise<Array<RedisServerDocker>> {
|
||||
const numberOfMasters = dockersConfig.numberOfMasters ?? 2,
|
||||
slotsPerNode = Math.floor(SLOTS / numberOfMasters),
|
||||
spawnPromises: Array<ReturnType<typeof spawnRedisClusterNodeDockers>> = [];
|
||||
for (let i = 0; i < numberOfMasters; i++) {
|
||||
const fromSlot = i * slotsPerNode,
|
||||
toSlot = i === numberOfMasters - 1 ? SLOTS : fromSlot + slotsPerNode;
|
||||
spawnPromises.push(
|
||||
spawnRedisClusterNodeDockers(
|
||||
dockersConfig,
|
||||
serverArguments,
|
||||
fromSlot,
|
||||
toSlot
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
const nodes = (await Promise.all(spawnPromises)).flat(),
|
||||
meetPromises: Array<Promise<unknown>> = [];
|
||||
for (let i = 1; i < nodes.length; i++) {
|
||||
meetPromises.push(
|
||||
nodes[i].client.clusterMeet('127.0.0.1', nodes[0].docker.port)
|
||||
);
|
||||
}
|
||||
|
||||
await Promise.all(meetPromises);
|
||||
|
||||
await Promise.all(
|
||||
nodes.map(async ({ client }) => {
|
||||
while (totalNodes(await client.clusterSlots()) !== nodes.length) {
|
||||
await promiseTimeout(50);
|
||||
}
|
||||
|
||||
return client.disconnect();
|
||||
})
|
||||
const numberOfMasters = dockersConfig.numberOfMasters ?? 2,
|
||||
slotsPerNode = Math.floor(SLOTS / numberOfMasters),
|
||||
spawnPromises: Array<ReturnType<typeof spawnRedisClusterNodeDockers>> = [];
|
||||
for (let i = 0; i < numberOfMasters; i++) {
|
||||
const fromSlot = i * slotsPerNode,
|
||||
toSlot = i === numberOfMasters - 1 ? SLOTS : fromSlot + slotsPerNode;
|
||||
spawnPromises.push(
|
||||
spawnRedisClusterNodeDockers(
|
||||
dockersConfig,
|
||||
serverArguments,
|
||||
fromSlot,
|
||||
toSlot
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
return nodes.map(({ docker }) => docker);
|
||||
const nodes = (await Promise.all(spawnPromises)).flat(),
|
||||
meetPromises: Array<Promise<unknown>> = [];
|
||||
for (let i = 1; i < nodes.length; i++) {
|
||||
meetPromises.push(
|
||||
nodes[i].client.clusterMeet('127.0.0.1', nodes[0].docker.port)
|
||||
);
|
||||
}
|
||||
|
||||
await Promise.all(meetPromises);
|
||||
|
||||
await Promise.all(
|
||||
nodes.map(async ({ client }) => {
|
||||
while (
|
||||
totalNodes(await client.clusterSlots()) !== nodes.length ||
|
||||
!(await client.sendCommand<string>(['CLUSTER', 'INFO'])).startsWith('cluster_state:ok') // TODO
|
||||
) {
|
||||
await setTimeout(50);
|
||||
}
|
||||
|
||||
client.destroy();
|
||||
})
|
||||
);
|
||||
|
||||
return nodes.map(({ docker }) => docker);
|
||||
}
|
||||
|
||||
function totalNodes(slots: ClusterSlotsReply) {
|
||||
let total = slots.length;
|
||||
for (const slot of slots) {
|
||||
total += slot.replicas.length;
|
||||
}
|
||||
// TODO: type ClusterSlotsReply
|
||||
function totalNodes(slots: any) {
|
||||
let total = slots.length;
|
||||
for (const slot of slots) {
|
||||
total += slot.replicas.length;
|
||||
}
|
||||
|
||||
return total;
|
||||
return total;
|
||||
}
|
||||
|
||||
const RUNNING_CLUSTERS = new Map<Array<string>, ReturnType<typeof spawnRedisClusterDockers>>();
|
||||
|
||||
export function spawnRedisCluster(dockersConfig: RedisClusterDockersConfig, serverArguments: Array<string>): Promise<Array<RedisServerDocker>> {
|
||||
const runningCluster = RUNNING_CLUSTERS.get(serverArguments);
|
||||
if (runningCluster) {
|
||||
return runningCluster;
|
||||
}
|
||||
const runningCluster = RUNNING_CLUSTERS.get(serverArguments);
|
||||
if (runningCluster) {
|
||||
return runningCluster;
|
||||
}
|
||||
|
||||
const dockersPromise = spawnRedisClusterDockers(dockersConfig, serverArguments);
|
||||
RUNNING_CLUSTERS.set(serverArguments, dockersPromise);
|
||||
return dockersPromise;
|
||||
const dockersPromise = spawnRedisClusterDockers(dockersConfig, serverArguments);
|
||||
RUNNING_CLUSTERS.set(serverArguments, dockersPromise);
|
||||
return dockersPromise;
|
||||
}
|
||||
|
||||
after(() => {
|
||||
return Promise.all(
|
||||
[...RUNNING_CLUSTERS.values()].map(async dockersPromise => {
|
||||
return Promise.all(
|
||||
(await dockersPromise).map(({ dockerId }) => dockerRemove(dockerId))
|
||||
);
|
||||
})
|
||||
);
|
||||
return Promise.all(
|
||||
[...RUNNING_CLUSTERS.values()].map(async dockersPromise => {
|
||||
return Promise.all(
|
||||
(await dockersPromise).map(({ dockerId }) => dockerRemove(dockerId))
|
||||
);
|
||||
})
|
||||
);
|
||||
});
|
||||
|
Reference in New Issue
Block a user