From 6e4122df67ae77b2592a5df8b17adb33580c8d92 Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Tue, 6 Jan 2026 14:16:16 +0200 Subject: [PATCH 1/5] Add predefined scenarios endpoint with remove-add scenario implementation - Add POST /scenarios/predefined/:scenario endpoint with Zod validation - Create scenarios folder structure similar to default_interceptors - Implement remove-add scenario: - Picks a random node from existing proxies - Adds a new node with auto-incremented port - Intercepts CLUSTER SLOTS to exclude picked node and include new node - Sends SMIGRATING notification (slots about to migrate) - Sends SMIGRATED notification (slots migrated to new node) - Add helper functions for scenario operations: - buildSMigratingNotification() - RESP3 SMIGRATING notification - buildSMigratedNotification() - RESP3 SMIGRATED notification - createCustomClusterSlotsInterceptor() - custom cluster slots response - getSlotRangesForProxy() - calculate slot ranges for a proxy - addNode() - add new proxy node - sendToAllClients() - broadcast to all clients - pickRandom() - random element selection - findNextAvailablePort() - find next available port - Add bar scenario skeleton for future implementation --- src/app.ts | 11 +++ src/scenarios/bar.ts | 18 ++++ src/scenarios/helpers.ts | 143 ++++++++++++++++++++++++++++++ src/scenarios/index.ts | 27 ++++++ src/scenarios/removeSrcAddDest.ts | 121 +++++++++++++++++++++++++ src/scenarios/sequence-gen.ts | 4 + src/util.ts | 4 + 7 files changed, 328 insertions(+) create mode 100644 src/scenarios/bar.ts create mode 100644 src/scenarios/helpers.ts create mode 100644 src/scenarios/index.ts create mode 100644 src/scenarios/removeSrcAddDest.ts create mode 100644 src/scenarios/sequence-gen.ts diff --git a/src/app.ts b/src/app.ts index e6d751f..0efcbf0 100644 --- a/src/app.ts +++ b/src/app.ts @@ -12,6 +12,7 @@ import { } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy.ts"; import applyDefaultInterceptors from "./default_interceptors/index.ts"; import ProxyStore, { makeId } from "./proxy-store.ts"; +import applyPredefinedScenario from "./scenarios/index.ts"; import { connectionIdsQuerySchema, type ExtendedProxyConfig, @@ -20,6 +21,7 @@ import { interceptorSchema, paramSchema, parseBuffer, + predefinedScenarioParamSchema, proxyConfigSchema, scenarioSchema, } from "./util.ts"; @@ -173,6 +175,15 @@ export function createApp(testConfig?: ExtendedProxyConfig) { return c.json({ success: true, totalResponses: responses.length }); }); + app.post( + "/scenarios/predefined/:scenario", + zValidator("param", predefinedScenarioParamSchema), + async (c) => { + const { scenario } = c.req.valid("param"); + return await applyPredefinedScenario(scenario, c, proxyStore, config); + }, + ); + app.post("/interceptors", zValidator("json", interceptorSchema), async (c) => { const { name, match, response, encoding } = c.req.valid("json"); diff --git a/src/scenarios/bar.ts b/src/scenarios/bar.ts new file mode 100644 index 0000000..b29f9bf --- /dev/null +++ b/src/scenarios/bar.ts @@ -0,0 +1,18 @@ +import type { Context } from "hono"; +import type ProxyStore from "../proxy-store"; +import type { ExtendedProxyConfig } from "../util"; + +export default async function barScenario( + c: Context, + _proxyStore: ProxyStore, + _config: ExtendedProxyConfig, +) { + // TODO: Implement bar scenario logic + // Example: Add interceptors, configure proxies, etc. + + return c.json({ + success: true, + scenario: "bar", + message: "Bar scenario executed", + }); +} diff --git a/src/scenarios/helpers.ts b/src/scenarios/helpers.ts new file mode 100644 index 0000000..1a24418 --- /dev/null +++ b/src/scenarios/helpers.ts @@ -0,0 +1,143 @@ +import { + type InterceptorDescription, + type InterceptorState, + type Next, + type ProxyConfig, + RedisProxy, + type SendResult, +} from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy"; +import type ProxyStore from "../proxy-store"; +import { makeId } from "../proxy-store"; + +/** + * Starts a new proxy with the given configuration + */ +export function startNewProxy(config: ProxyConfig): RedisProxy { + const proxy = new RedisProxy(config); + proxy.start().catch(console.error); + return proxy; +} + +/** + * Adds a new node to the proxy store + */ +export function addNode( + proxyStore: ProxyStore, + config: ProxyConfig, +): { nodeId: string; proxy: RedisProxy } { + const nodeId = makeId(config.targetHost, config.targetPort, config.listenPort); + const proxy = startNewProxy(config); + proxyStore.add(nodeId, proxy); + return { nodeId, proxy }; +} + +/** + * Sends a buffer to all clients across all proxies + */ +export function sendToAllClients(proxyStore: ProxyStore, buffer: Buffer): SendResult[] { + const results: SendResult[] = []; + for (const proxy of proxyStore.proxies) { + results.push(...proxy.sendToAllClients(buffer)); + } + return results; +} + +/** + * Creates a cluster slots interceptor that returns a custom list of proxies + */ +export function createCustomClusterSlotsInterceptor( + proxiesToInclude: RedisProxy[], +): InterceptorDescription { + return { + name: "cluster-simulation-interceptor", + fn: async (data: Buffer, next: Next, state: InterceptorState) => { + state.invokeCount++; + + if (data.toString().toLowerCase() !== "*2\r\n$7\r\ncluster\r\n$5\r\nslots\r\n") { + return next(data); + } + + state.matchCount++; + + const slotLength = Math.floor(16384 / proxiesToInclude.length); + + let current = -1; + const mapping = proxiesToInclude.map((proxy, i) => { + const from = current + 1; + const to = i === proxiesToInclude.length - 1 ? 16383 : current + slotLength; + current = to; + const id = `proxy-id-${proxy.config.listenPort}`; + return `*3\r\n:${from}\r\n:${to}\r\n*3\r\n$${proxy.config.listenHost.length}\r\n${proxy.config.listenHost}\r\n:${proxy.config.listenPort}\r\n$${id.length}\r\n${id}\r\n`; + }); + + const response = `*${proxiesToInclude.length}\r\n${mapping.join("")}`; + return Buffer.from(response); + }, + }; +} + +/** + * Picks a random element from an array + */ +export function pickRandom(array: T[]): T | undefined { + if (array.length === 0) return undefined; + return array[Math.floor(Math.random() * array.length)]; +} + +/** + * Finds the next available port by incrementing from the highest existing port + */ +export function findNextAvailablePort(proxies: RedisProxy[]): number { + const ports = proxies.map((p) => p.config.listenPort); + return Math.max(...ports) + 1; +} + +/** + * Builds an SMIGRATING notification in RESP3 format + * This notifies clients that slots are about to be migrated + */ +export function buildSMigratingNotification(slotRanges: string, seqId: number = 1): Buffer { + const response = `>3\r\n+SMIGRATING\r\n:${seqId}\r\n+${slotRanges}\r\n`; + return Buffer.from(response); +} + +/** + * Builds an SMIGRATED notification in RESP3 format + * This notifies clients that slots have been migrated to different nodes + */ +export function buildSMigratedNotification( + movedSlotsByDestination: Array<{ + targetNode: { host: string; port: number }; + slotRanges: string; // e.g., "0-5460" or "0-100,200-300,500" + }>, + seqId: number = 1, +): Buffer { + if (movedSlotsByDestination.length === 0) { + throw new Error("No slots to migrate"); + } + + const entries = movedSlotsByDestination.map(({ targetNode, slotRanges }) => { + const hostPort = `${targetNode.host}:${targetNode.port}`; + return `*2\r\n+${hostPort}\r\n+${slotRanges}\r\n`; + }); + + const response = `>3\r\n+SMIGRATED\r\n:${seqId}\r\n*${movedSlotsByDestination.length}\r\n${entries.join("")}`; + + return Buffer.from(response); +} + +/** + * Gets the slot ranges assigned to a specific proxy based on cluster slot distribution + */ +export function getSlotRangesForProxy(proxy: RedisProxy, allProxies: RedisProxy[]): string { + const proxyIndex = allProxies.indexOf(proxy); + if (proxyIndex === -1) { + throw new Error("Proxy not found in the list"); + } + + const slotLength = Math.floor(16384 / allProxies.length); + const from = proxyIndex * slotLength; + const to = proxyIndex === allProxies.length - 1 ? 16383 : from + slotLength - 1; + + return `${from}-${to}`; +} diff --git a/src/scenarios/index.ts b/src/scenarios/index.ts new file mode 100644 index 0000000..ca8e309 --- /dev/null +++ b/src/scenarios/index.ts @@ -0,0 +1,27 @@ +import type { Context } from "hono"; +import type { z } from "zod"; +import type ProxyStore from "../proxy-store"; +import type { ExtendedProxyConfig, predefinedScenarioParamSchema } from "../util"; +import barScenario from "./bar"; +import removeSrcAddDestScenario from "./removeSrcAddDest"; + +type PredefinedScenario = z.infer["scenario"]; + +export default async function applyPredefinedScenario( + scenario: PredefinedScenario, + c: Context, + proxyStore: ProxyStore, + config: ExtendedProxyConfig, +) { + switch (scenario) { + case "remove-add": + return await removeSrcAddDestScenario(c, proxyStore, config); + + case "bar": + return await barScenario(c, proxyStore, config); + + default: + // This should never happen due to Zod validation, but TypeScript requires it + return c.json({ success: false, error: "Unknown scenario" }, 400); + } +} diff --git a/src/scenarios/removeSrcAddDest.ts b/src/scenarios/removeSrcAddDest.ts new file mode 100644 index 0000000..a2b7570 --- /dev/null +++ b/src/scenarios/removeSrcAddDest.ts @@ -0,0 +1,121 @@ +import type { Context } from "hono"; +import type { ProxyConfig } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy"; +import type ProxyStore from "../proxy-store"; + +import { makeId } from "../proxy-store"; +import type { ExtendedProxyConfig } from "../util"; +import { + addNode, + buildSMigratedNotification, + buildSMigratingNotification, + createCustomClusterSlotsInterceptor, + findNextAvailablePort, + getSlotRangesForProxy, + pickRandom, +} from "./helpers"; +import { getNextSequenceId } from "./sequence-gen"; + +/** + * "REMOVED AND ADDED" Scenario: + * + * 1. Pick a random node from existing proxies + * 2. Add one more node + * 3. Intercept CLUSTER SLOTS to return all nodes + the new one - the randomly selected one + * 4. Send SMIGRATING notification to all clients (slots about to be migrated) + * 5. Send SMIGRATED notification to all clients (slots migrated from picked node to new node) + * 6. Kill the old node + */ +export default async function removeSrcAddDestScenario( + c: Context, + proxyStore: ProxyStore, + config: ExtendedProxyConfig, +) { + // Step 1: Pick a random node + const allProxies = proxyStore.proxies; + if (allProxies.length === 0) { + return c.json( + { + success: false, + error: "No proxies available to select from", + }, + 400, + ); + } + + const proxyToBeRemoved = pickRandom(allProxies); + if (!proxyToBeRemoved) { + return c.json( + { + success: false, + error: "Failed to select a random proxy", + }, + 400, + ); + } + + // Get the slot ranges for the randomly selected proxy before adding the new node + const slotRanges = getSlotRangesForProxy(proxyToBeRemoved, allProxies); + + // Step 2: Add one more node + const newPort = findNextAvailablePort(allProxies); + const newProxyConfig: ProxyConfig = { + ...config, + listenPort: newPort, + }; + + const { nodeId: newNodeId, proxy: newProxy } = addNode(proxyStore, newProxyConfig); + + // Step 3: Create list of proxies excluding the randomly selected one, and including the new one + const proxiesForClusterSlots = allProxies.filter((p) => p !== proxyToBeRemoved).concat(newProxy); + + // Add the custom cluster slots interceptor to all proxies + const clusterSlotsInterceptor = createCustomClusterSlotsInterceptor(proxiesForClusterSlots); + + for (const proxy of proxyStore.proxies) { + proxy.addGlobalInterceptor(clusterSlotsInterceptor); + } + + // Step 4: Send SMIGRATING notification to all clients + // This notifies clients that slots are about to be migrated + const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); + proxyToBeRemoved.sendToAllClients(sMigratingBuffer); + + // Step 5: Send SMIGRATED notification to all clients + // This notifies clients that slots from the picked node have moved to the new node + setTimeout(() => { + const sMigratedBuffer = buildSMigratedNotification( + [ + { + targetNode: { + host: newProxy.config.listenHost, + port: newProxy.config.listenPort, + }, + slotRanges, + }, + ], + getNextSequenceId(), + ); + proxyToBeRemoved.sendToAllClients(sMigratedBuffer); + + setTimeout(() => { + const { targetHost, targetPort, listenPort } = proxyToBeRemoved.config; + proxyStore.delete(makeId(targetHost, targetPort, listenPort)); + }, 2000); + }, 5000); + + return c.json({ + success: true, + scenario: "foo", + message: "Foo scenario started successfully", + details: { + excludedNode: `${proxyToBeRemoved.config.listenHost}:${proxyToBeRemoved.config.listenPort}`, + newNode: `${newProxy.config.listenHost}:${newProxy.config.listenPort}`, + newNodeId, + migratedSlots: slotRanges, + clusterSlotsProxies: proxiesForClusterSlots.map((p) => ({ + host: p.config.listenHost, + port: p.config.listenPort, + })), + }, + }); +} diff --git a/src/scenarios/sequence-gen.ts b/src/scenarios/sequence-gen.ts new file mode 100644 index 0000000..8d21a31 --- /dev/null +++ b/src/scenarios/sequence-gen.ts @@ -0,0 +1,4 @@ +let id = 1; +export function getNextSequenceId() { + return id++; +} diff --git a/src/util.ts b/src/util.ts index cabade3..6cb9ed1 100644 --- a/src/util.ts +++ b/src/util.ts @@ -36,6 +36,10 @@ export const interceptorSchema = z.object({ response: z.string(), }); +export const predefinedScenarioParamSchema = z.object({ + scenario: z.enum(["remove-add", "bar"]), +}); + export function parseBuffer(data: string, encoding: "base64" | "raw"): Buffer { switch (encoding) { case "base64": From 818d92f685dd7f785092dd39377e7b4895ebd411 Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Fri, 9 Jan 2026 13:44:07 +0200 Subject: [PATCH 2/5] Implement add, remove, shuffle scenarios --- src/scenarios/add.ts | 119 ++++++++++++++ src/scenarios/bar.ts | 18 --- src/scenarios/index.ts | 46 +++--- .../{removeSrcAddDest.ts => remove-add.ts} | 0 src/scenarios/remove.ts | 125 +++++++++++++++ src/scenarios/slot-shuffle.ts | 151 ++++++++++++++++++ src/util.ts | 2 +- 7 files changed, 424 insertions(+), 37 deletions(-) create mode 100644 src/scenarios/add.ts delete mode 100644 src/scenarios/bar.ts rename src/scenarios/{removeSrcAddDest.ts => remove-add.ts} (100%) create mode 100644 src/scenarios/remove.ts create mode 100644 src/scenarios/slot-shuffle.ts diff --git a/src/scenarios/add.ts b/src/scenarios/add.ts new file mode 100644 index 0000000..960bb2c --- /dev/null +++ b/src/scenarios/add.ts @@ -0,0 +1,119 @@ +import type { Context } from "hono"; +import type { ProxyConfig } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy"; +import type ProxyStore from "../proxy-store"; + +import type { ExtendedProxyConfig } from "../util"; +import { + addNode, + buildSMigratedNotification, + buildSMigratingNotification, + createCustomClusterSlotsInterceptor, + findNextAvailablePort, + getSlotRangesForProxy, +} from "./helpers"; +import { getNextSequenceId } from "./sequence-gen"; + +/** + * "ADD NODE" Scenario: + * + * 1. Add a new node to the cluster + * 2. Calculate new slot distribution (slots redistributed from existing nodes to new node) + * 3. Intercept CLUSTER SLOTS to return all nodes (existing + new) with updated slot ranges + * 4. Send SMIGRATING notifications from nodes that will lose slots + * 5. Send SMIGRATED notifications indicating slots moved to the new node + * 6. All nodes remain active + */ +export default async function addNodeScenario( + c: Context, + proxyStore: ProxyStore, + config: ExtendedProxyConfig, +) { + const allProxies = proxyStore.proxies; + if (allProxies.length === 0) { + return c.json( + { + success: false, + error: "No proxies available", + }, + 400, + ); + } + + // Step 1: Get current slot distribution before adding new node + const oldSlotDistribution = allProxies.map((proxy) => ({ + proxy, + slotRanges: getSlotRangesForProxy(proxy, allProxies), + })); + + // Step 2: Add a new node + const newPort = findNextAvailablePort(allProxies); + const newProxyConfig: ProxyConfig = { + ...config, + listenPort: newPort, + }; + + const { nodeId: newNodeId, proxy: newProxy } = addNode(proxyStore, newProxyConfig); + + // Step 3: Calculate new slot distribution with the new node included + const allProxiesWithNew = [...allProxies, newProxy]; + const newSlotDistribution = allProxiesWithNew.map((proxy) => ({ + proxy, + slotRanges: getSlotRangesForProxy(proxy, allProxiesWithNew), + })); + + // Step 4: Intercept CLUSTER SLOTS to return all nodes with new distribution + const clusterSlotsInterceptor = createCustomClusterSlotsInterceptor(allProxiesWithNew); + + for (const proxy of proxyStore.proxies) { + proxy.addGlobalInterceptor(clusterSlotsInterceptor); + } + + // Step 5: Send SMIGRATING notifications from nodes that will lose slots + // Each existing node loses some slots to make room for the new node + for (const { proxy, slotRanges } of oldSlotDistribution) { + const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); + proxy.sendToAllClients(sMigratingBuffer); + } + + // Step 6: Send SMIGRATED notifications after a delay + // This indicates where slots have moved (to the new node) + setTimeout(() => { + const newNodeSlotRanges = getSlotRangesForProxy(newProxy, allProxiesWithNew); + + // Send SMIGRATED from each old node indicating their slots moved to new node + for (const { proxy } of oldSlotDistribution) { + const sMigratedBuffer = buildSMigratedNotification( + [ + { + targetNode: { + host: newProxy.config.listenHost, + port: newProxy.config.listenPort, + }, + slotRanges: newNodeSlotRanges, + }, + ], + getNextSequenceId(), + ); + proxy.sendToAllClients(sMigratedBuffer); + } + }, 5000); + + return c.json({ + success: true, + scenario: "add", + message: "Add node scenario started successfully", + details: { + newNode: `${newProxy.config.listenHost}:${newProxy.config.listenPort}`, + newNodeId, + newNodeSlots: getSlotRangesForProxy(newProxy, allProxiesWithNew), + oldDistribution: oldSlotDistribution.map(({ proxy, slotRanges }) => ({ + node: `${proxy.config.listenHost}:${proxy.config.listenPort}`, + slots: slotRanges, + })), + newDistribution: newSlotDistribution.map(({ proxy, slotRanges }) => ({ + node: `${proxy.config.listenHost}:${proxy.config.listenPort}`, + slots: slotRanges, + })), + }, + }); +} diff --git a/src/scenarios/bar.ts b/src/scenarios/bar.ts deleted file mode 100644 index b29f9bf..0000000 --- a/src/scenarios/bar.ts +++ /dev/null @@ -1,18 +0,0 @@ -import type { Context } from "hono"; -import type ProxyStore from "../proxy-store"; -import type { ExtendedProxyConfig } from "../util"; - -export default async function barScenario( - c: Context, - _proxyStore: ProxyStore, - _config: ExtendedProxyConfig, -) { - // TODO: Implement bar scenario logic - // Example: Add interceptors, configure proxies, etc. - - return c.json({ - success: true, - scenario: "bar", - message: "Bar scenario executed", - }); -} diff --git a/src/scenarios/index.ts b/src/scenarios/index.ts index ca8e309..db5e4c2 100644 --- a/src/scenarios/index.ts +++ b/src/scenarios/index.ts @@ -1,27 +1,37 @@ import type { Context } from "hono"; import type { z } from "zod"; import type ProxyStore from "../proxy-store"; -import type { ExtendedProxyConfig, predefinedScenarioParamSchema } from "../util"; -import barScenario from "./bar"; -import removeSrcAddDestScenario from "./removeSrcAddDest"; +import type { + ExtendedProxyConfig, + predefinedScenarioParamSchema, +} from "../util"; +import addNodeScenario from "./add"; +import removeSrcAddDestScenario from "./remove-add"; +import removeNodeScenario from "./remove"; +import slotShuffleScenario from "./slot-shuffle"; -type PredefinedScenario = z.infer["scenario"]; +type PredefinedScenario = z.infer< + typeof predefinedScenarioParamSchema +>["scenario"]; export default async function applyPredefinedScenario( - scenario: PredefinedScenario, - c: Context, - proxyStore: ProxyStore, - config: ExtendedProxyConfig, + scenario: PredefinedScenario, + c: Context, + proxyStore: ProxyStore, + config: ExtendedProxyConfig, ) { - switch (scenario) { - case "remove-add": - return await removeSrcAddDestScenario(c, proxyStore, config); + switch (scenario) { + case "remove-add": + return await removeSrcAddDestScenario(c, proxyStore, config); + case "add": + return await addNodeScenario(c, proxyStore, config); + case "remove": + return await removeNodeScenario(c, proxyStore, config); + case "slot-shuffle": + return await slotShuffleScenario(c, proxyStore, config); - case "bar": - return await barScenario(c, proxyStore, config); - - default: - // This should never happen due to Zod validation, but TypeScript requires it - return c.json({ success: false, error: "Unknown scenario" }, 400); - } + default: + // This should never happen due to Zod validation, but TypeScript requires it + return c.json({ success: false, error: "Unknown scenario" }, 400); + } } diff --git a/src/scenarios/removeSrcAddDest.ts b/src/scenarios/remove-add.ts similarity index 100% rename from src/scenarios/removeSrcAddDest.ts rename to src/scenarios/remove-add.ts diff --git a/src/scenarios/remove.ts b/src/scenarios/remove.ts new file mode 100644 index 0000000..da02c1e --- /dev/null +++ b/src/scenarios/remove.ts @@ -0,0 +1,125 @@ +import type { Context } from "hono"; +import type ProxyStore from "../proxy-store"; + +import { makeId } from "../proxy-store"; +import type { ExtendedProxyConfig } from "../util"; +import { + buildSMigratedNotification, + buildSMigratingNotification, + createCustomClusterSlotsInterceptor, + getSlotRangesForProxy, + pickRandom, +} from "./helpers"; +import { getNextSequenceId } from "./sequence-gen"; + +/** + * "REMOVE NODE" Scenario: + * + * 1. Pick a random node from existing proxies to remove + * 2. Get the slot ranges currently assigned to that node + * 3. Calculate new slot distribution among remaining nodes + * 4. Intercept CLUSTER SLOTS to return all nodes except the one being removed + * 5. Send SMIGRATING notification from the node being removed + * 6. Send SMIGRATED notifications indicating where slots moved (to remaining nodes) + * 7. Kill/delete the removed node after notifications + */ +export default async function removeNodeScenario( + c: Context, + proxyStore: ProxyStore, + _config: ExtendedProxyConfig, +) { + // Step 1: Pick a random node to remove + const allProxies = proxyStore.proxies; + if (allProxies.length === 0) { + return c.json( + { + success: false, + error: "No proxies available to select from", + }, + 400, + ); + } + + if (allProxies.length === 1) { + return c.json( + { + success: false, + error: "Cannot remove the last remaining node", + }, + 400, + ); + } + + const proxyToBeRemoved = pickRandom(allProxies); + if (!proxyToBeRemoved) { + return c.json( + { + success: false, + error: "Failed to select a random proxy", + }, + 400, + ); + } + + // Step 2: Get the slot ranges for the node being removed + const removedNodeSlotRanges = getSlotRangesForProxy(proxyToBeRemoved, allProxies); + + // Step 3: Create list of remaining proxies (excluding the one being removed) + const remainingProxies = allProxies.filter((p) => p !== proxyToBeRemoved); + + // Calculate new slot distribution for remaining nodes + const newSlotDistribution = remainingProxies.map((proxy) => ({ + proxy, + slotRanges: getSlotRangesForProxy(proxy, remainingProxies), + })); + + // Step 4: Intercept CLUSTER SLOTS to return only remaining nodes + const clusterSlotsInterceptor = createCustomClusterSlotsInterceptor(remainingProxies); + + for (const proxy of proxyStore.proxies) { + proxy.addGlobalInterceptor(clusterSlotsInterceptor); + } + + // Step 5: Send SMIGRATING notification from the node being removed + const sMigratingBuffer = buildSMigratingNotification(removedNodeSlotRanges, getNextSequenceId()); + proxyToBeRemoved.sendToAllClients(sMigratingBuffer); + + // Step 6: Send SMIGRATED notifications after a delay + // This indicates where slots from the removed node have moved (distributed to remaining nodes) + setTimeout(() => { + // Build SMIGRATED notification showing slots redistributed to remaining nodes + const migratedSlots = newSlotDistribution.map(({ proxy, slotRanges }) => ({ + targetNode: { + host: proxy.config.listenHost, + port: proxy.config.listenPort, + }, + slotRanges, + })); + + const sMigratedBuffer = buildSMigratedNotification(migratedSlots, getNextSequenceId()); + proxyToBeRemoved.sendToAllClients(sMigratedBuffer); + + // Step 7: Kill the removed node after another delay + setTimeout(() => { + const { targetHost, targetPort, listenPort } = proxyToBeRemoved.config; + proxyStore.delete(makeId(targetHost, targetPort, listenPort)); + }, 2000); + }, 5000); + + return c.json({ + success: true, + scenario: "remove", + message: "Remove node scenario started successfully", + details: { + removedNode: `${proxyToBeRemoved.config.listenHost}:${proxyToBeRemoved.config.listenPort}`, + removedNodeSlots: removedNodeSlotRanges, + remainingNodes: remainingProxies.map((p) => ({ + node: `${p.config.listenHost}:${p.config.listenPort}`, + })), + newDistribution: newSlotDistribution.map(({ proxy, slotRanges }) => ({ + node: `${proxy.config.listenHost}:${proxy.config.listenPort}`, + slots: slotRanges, + })), + }, + }); +} diff --git a/src/scenarios/slot-shuffle.ts b/src/scenarios/slot-shuffle.ts new file mode 100644 index 0000000..60d7189 --- /dev/null +++ b/src/scenarios/slot-shuffle.ts @@ -0,0 +1,151 @@ +import type { Context } from "hono"; +import type { RedisProxy } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy"; +import type ProxyStore from "../proxy-store"; + +import type { ExtendedProxyConfig } from "../util"; +import { + buildSMigratedNotification, + buildSMigratingNotification, + createCustomClusterSlotsInterceptor, + getSlotRangesForProxy, +} from "./helpers"; +import { getNextSequenceId } from "./sequence-gen"; + +/** + * "SLOT SHUFFLE" Scenario: + * + * 1. Keep all existing nodes (no add/remove) + * 2. Randomly shuffle slot assignments among existing nodes + * 3. Intercept CLUSTER SLOTS to return the new shuffled slot distribution + * 4. Send SMIGRATING notifications from nodes losing slots + * 5. Send SMIGRATED notifications indicating the new slot ownership + * 6. All nodes remain active with new slot assignments + */ +export default async function slotShuffleScenario( + c: Context, + proxyStore: ProxyStore, + _config: ExtendedProxyConfig, +) { + const allProxies = proxyStore.proxies; + if (allProxies.length === 0) { + return c.json( + { + success: false, + error: "No proxies available", + }, + 400, + ); + } + + if (allProxies.length === 1) { + return c.json( + { + success: false, + error: "Cannot shuffle slots with only one node", + }, + 400, + ); + } + + // Step 1: Get current slot distribution + const oldSlotDistribution = allProxies.map((proxy) => ({ + proxy, + slotRanges: getSlotRangesForProxy(proxy, allProxies), + })); + + // Step 2: Shuffle the proxies array to create a new slot distribution + // We'll create a shuffled copy of the proxies array + const shuffledProxies = [...allProxies]; + + // Fisher-Yates shuffle algorithm + for (let i = shuffledProxies.length - 1; i > 0; i--) { + const j = Math.floor(Math.random() * (i + 1)); + const temp = shuffledProxies[i]; + const jProxy = shuffledProxies[j]; + if (temp && jProxy) { + shuffledProxies[i] = jProxy; + shuffledProxies[j] = temp; + } + } + + // Step 3: Calculate new slot distribution based on shuffled order + // We'll use the shuffled order to reassign slots + const newSlotDistribution = shuffledProxies.map((proxy, index) => { + const slotLength = Math.floor(16384 / shuffledProxies.length); + const from = index * slotLength; + const to = index === shuffledProxies.length - 1 ? 16383 : from + slotLength - 1; + return { + proxy, + slotRanges: `${from}-${to}`, + }; + }); + + // Step 4: Create a custom interceptor that returns the shuffled slot distribution + // We need to create a custom interceptor because the slots don't match the original order + const customShuffledInterceptor = { + name: "cluster-simulation-interceptor", + fn: async (data: Buffer, next: any, state: any) => { + state.invokeCount++; + + if (data.toString().toLowerCase() !== "*2\r\n$7\r\ncluster\r\n$5\r\nslots\r\n") { + return next(data); + } + + state.matchCount++; + + const mapping = newSlotDistribution.map(({ proxy, slotRanges }) => { + const [from, to] = slotRanges.split("-").map(Number); + const id = `proxy-id-${proxy.config.listenPort}`; + return `*3\r\n:${from}\r\n:${to}\r\n*3\r\n$${proxy.config.listenHost.length}\r\n${proxy.config.listenHost}\r\n:${proxy.config.listenPort}\r\n$${id.length}\r\n${id}\r\n`; + }); + + const response = `*${newSlotDistribution.length}\r\n${mapping.join("")}`; + return Buffer.from(response); + }, + }; + + for (const proxy of proxyStore.proxies) { + proxy.addGlobalInterceptor(customShuffledInterceptor); + } + + // Step 5: Send SMIGRATING notifications from all nodes + // All nodes are potentially losing and gaining slots + for (const { proxy, slotRanges } of oldSlotDistribution) { + const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); + proxy.sendToAllClients(sMigratingBuffer); + } + + // Step 6: Send SMIGRATED notifications after a delay + // This indicates the new slot ownership across all nodes + setTimeout(() => { + // Each node sends SMIGRATED showing where all slots now live + for (const { proxy: sourceProxy } of oldSlotDistribution) { + const migratedSlots = newSlotDistribution.map(({ proxy, slotRanges }) => ({ + targetNode: { + host: proxy.config.listenHost, + port: proxy.config.listenPort, + }, + slotRanges, + })); + + const sMigratedBuffer = buildSMigratedNotification(migratedSlots, getNextSequenceId()); + sourceProxy.sendToAllClients(sMigratedBuffer); + } + }, 5000); + + return c.json({ + success: true, + scenario: "slot-shuffle", + message: "Slot shuffle scenario started successfully", + details: { + oldDistribution: oldSlotDistribution.map(({ proxy, slotRanges }) => ({ + node: `${proxy.config.listenHost}:${proxy.config.listenPort}`, + slots: slotRanges, + })), + newDistribution: newSlotDistribution.map(({ proxy, slotRanges }) => ({ + node: `${proxy.config.listenHost}:${proxy.config.listenPort}`, + slots: slotRanges, + })), + }, + }); +} diff --git a/src/util.ts b/src/util.ts index 6cb9ed1..87e5a24 100644 --- a/src/util.ts +++ b/src/util.ts @@ -37,7 +37,7 @@ export const interceptorSchema = z.object({ }); export const predefinedScenarioParamSchema = z.object({ - scenario: z.enum(["remove-add", "bar"]), + scenario: z.enum(["remove-add", "remove", "add", "slot-shuffle"]), }); export function parseBuffer(data: string, encoding: "base64" | "raw"): Buffer { From 844b132100206109b7f246bf4184bc5b42c33bf1 Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Fri, 23 Jan 2026 09:21:55 +0200 Subject: [PATCH 3/5] feat(actions): add Fault Injector-compatible action execution API Implement action-based API matching the Fault Injector interface for executing Redis cluster operations with slot migration simulation. New API endpoints: - POST /action - Submit an action for async execution - GET /action/:action_id - Poll action status and result - GET /slot-migrate - List triggers for slot migration effects Implemented slot migration effects: - remove-add: Remove one node and add a new one - remove: Remove a node and redistribute slots - add: Add a new node and rebalance slots - slot-shuffle: Redistribute slots across existing nodes Each effect sends SMIGRATING/SMIGRATED push notifications to connected clients and updates CLUSTER SLOTS interceptors accordingly. Removed deprecated /scenarios endpoints in favor of the new action API. --- src/actions/index.ts | 283 ++++++++++++++++++++++++++++++++++ src/app.ts | 184 +++++++++++++++++----- src/scenarios/slot-shuffle.ts | 1 - src/util.ts | 79 +++++++++- 4 files changed, 502 insertions(+), 45 deletions(-) create mode 100644 src/actions/index.ts diff --git a/src/actions/index.ts b/src/actions/index.ts new file mode 100644 index 0000000..9b46467 --- /dev/null +++ b/src/actions/index.ts @@ -0,0 +1,283 @@ +import type { ProxyConfig } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy"; +import type ProxyStore from "../proxy-store"; +import { makeId } from "../proxy-store"; +import type { ActionType, ExtendedProxyConfig } from "../util"; +import { + addNode, + buildSMigratedNotification, + buildSMigratingNotification, + createCustomClusterSlotsInterceptor, + findNextAvailablePort, + getSlotRangesForProxy, + pickRandom, +} from "../scenarios/helpers"; +import { getNextSequenceId } from "../scenarios/sequence-gen"; + +// Effect types matching Python MigrateEffect enum +export type SlotMigrateEffect = "remove-add" | "remove" | "add" | "slot-shuffle"; + +export interface SlotMigrateParams { + effect: SlotMigrateEffect; + variant?: string; + source_node?: number; + target_node?: number; +} + +export interface ActionExecutionResult { + status: "success" | "failed"; + error?: string | null; +} + +const delay = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); + +/** + * Execute an action based on its type and parameters + */ +export async function executeAction( + actionType: ActionType, + parameters: Record, + proxyStore: ProxyStore, + config: ExtendedProxyConfig, +): Promise { + switch (actionType) { + case "slot_migrate": + return executeSlotMigrate(parameters as unknown as SlotMigrateParams, proxyStore, config); + default: + return { status: "success" }; + } +} + +/** + * Execute slot_migrate action - handles all scenario effects + */ +async function executeSlotMigrate( + params: SlotMigrateParams, + proxyStore: ProxyStore, + config: ExtendedProxyConfig, +): Promise { + const { effect } = params; + + if (!effect) { + return { status: "failed", error: "Missing required parameter: effect" }; + } + + const validEffects: SlotMigrateEffect[] = ["remove-add", "remove", "add", "slot-shuffle"]; + if (!validEffects.includes(effect)) { + return { status: "failed", error: `Invalid effect: ${effect}. Must be one of: ${validEffects.join(", ")}` }; + } + + try { + switch (effect) { + case "remove-add": + await executeRemoveAddEffect(proxyStore, config); + break; + case "remove": + await executeRemoveEffect(proxyStore, config); + break; + case "add": + await executeAddEffect(proxyStore, config); + break; + case "slot-shuffle": + await executeSlotShuffleEffect(proxyStore, config); + break; + } + return { status: "success" }; + } catch (error) { + return { status: "failed", error: error instanceof Error ? error.message : String(error) }; + } +} + +async function executeRemoveAddEffect(proxyStore: ProxyStore, config: ExtendedProxyConfig): Promise { + const allProxies = proxyStore.proxies; + if (allProxies.length === 0) { + throw new Error("No proxies available to select from"); + } + + const proxyToBeRemoved = pickRandom(allProxies); + if (!proxyToBeRemoved) { + throw new Error("Failed to select a random proxy"); + } + + const slotRanges = getSlotRangesForProxy(proxyToBeRemoved, allProxies); + const newPort = findNextAvailablePort(allProxies); + const newProxyConfig: ProxyConfig = { ...config, listenPort: newPort }; + const { proxy: newProxy } = addNode(proxyStore, newProxyConfig); + + const proxiesForClusterSlots = allProxies.filter((p) => p !== proxyToBeRemoved).concat(newProxy); + const clusterSlotsInterceptor = createCustomClusterSlotsInterceptor(proxiesForClusterSlots); + + for (const proxy of proxyStore.proxies) { + proxy.addGlobalInterceptor(clusterSlotsInterceptor); + } + + const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); + proxyToBeRemoved.sendToAllClients(sMigratingBuffer); + + await delay(5000); + + const sMigratedBuffer = buildSMigratedNotification( + [{ targetNode: { host: newProxy.config.listenHost, port: newProxy.config.listenPort }, slotRanges }], + getNextSequenceId(), + ); + proxyToBeRemoved.sendToAllClients(sMigratedBuffer); + + await delay(2000); + + const { targetHost, targetPort, listenPort } = proxyToBeRemoved.config; + await proxyStore.delete(makeId(targetHost, targetPort, listenPort)); +} + +async function executeRemoveEffect(proxyStore: ProxyStore, _config: ExtendedProxyConfig): Promise { + const allProxies = proxyStore.proxies; + if (allProxies.length === 0) { + throw new Error("No proxies available to select from"); + } + if (allProxies.length === 1) { + throw new Error("Cannot remove the last remaining node"); + } + + const proxyToBeRemoved = pickRandom(allProxies); + if (!proxyToBeRemoved) { + throw new Error("Failed to select a random proxy"); + } + + const removedNodeSlotRanges = getSlotRangesForProxy(proxyToBeRemoved, allProxies); + const remainingProxies = allProxies.filter((p) => p !== proxyToBeRemoved); + + const newSlotDistribution = remainingProxies.map((proxy) => ({ + proxy, + slotRanges: getSlotRangesForProxy(proxy, remainingProxies), + })); + + const clusterSlotsInterceptor = createCustomClusterSlotsInterceptor(remainingProxies); + for (const proxy of proxyStore.proxies) { + proxy.addGlobalInterceptor(clusterSlotsInterceptor); + } + + const sMigratingBuffer = buildSMigratingNotification(removedNodeSlotRanges, getNextSequenceId()); + proxyToBeRemoved.sendToAllClients(sMigratingBuffer); + + await delay(5000); + + const migratedSlots = newSlotDistribution.map(({ proxy, slotRanges }) => ({ + targetNode: { host: proxy.config.listenHost, port: proxy.config.listenPort }, + slotRanges, + })); + const sMigratedBuffer = buildSMigratedNotification(migratedSlots, getNextSequenceId()); + proxyToBeRemoved.sendToAllClients(sMigratedBuffer); + + await delay(2000); + + const { targetHost, targetPort, listenPort } = proxyToBeRemoved.config; + await proxyStore.delete(makeId(targetHost, targetPort, listenPort)); +} + +async function executeAddEffect(proxyStore: ProxyStore, config: ExtendedProxyConfig): Promise { + const allProxies = proxyStore.proxies; + if (allProxies.length === 0) { + throw new Error("No proxies available"); + } + + const oldSlotDistribution = allProxies.map((proxy) => ({ + proxy, + slotRanges: getSlotRangesForProxy(proxy, allProxies), + })); + + const newPort = findNextAvailablePort(allProxies); + const newProxyConfig: ProxyConfig = { ...config, listenPort: newPort }; + const { proxy: newProxy } = addNode(proxyStore, newProxyConfig); + + const allProxiesWithNew = [...allProxies, newProxy]; + const clusterSlotsInterceptor = createCustomClusterSlotsInterceptor(allProxiesWithNew); + + for (const proxy of proxyStore.proxies) { + proxy.addGlobalInterceptor(clusterSlotsInterceptor); + } + + for (const { proxy, slotRanges } of oldSlotDistribution) { + const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); + proxy.sendToAllClients(sMigratingBuffer); + } + + await delay(5000); + + const newNodeSlotRanges = getSlotRangesForProxy(newProxy, allProxiesWithNew); + for (const { proxy } of oldSlotDistribution) { + const sMigratedBuffer = buildSMigratedNotification( + [{ targetNode: { host: newProxy.config.listenHost, port: newProxy.config.listenPort }, slotRanges: newNodeSlotRanges }], + getNextSequenceId(), + ); + proxy.sendToAllClients(sMigratedBuffer); + } +} + +async function executeSlotShuffleEffect(proxyStore: ProxyStore, _config: ExtendedProxyConfig): Promise { + const allProxies = proxyStore.proxies; + if (allProxies.length === 0) { + throw new Error("No proxies available"); + } + if (allProxies.length === 1) { + throw new Error("Cannot shuffle slots with only one node"); + } + + const oldSlotDistribution = allProxies.map((proxy) => ({ + proxy, + slotRanges: getSlotRangesForProxy(proxy, allProxies), + })); + + // Fisher-Yates shuffle + const shuffledProxies = [...allProxies]; + for (let i = shuffledProxies.length - 1; i > 0; i--) { + const j = Math.floor(Math.random() * (i + 1)); + const temp = shuffledProxies[i]; + const jProxy = shuffledProxies[j]; + if (temp && jProxy) { + shuffledProxies[i] = jProxy; + shuffledProxies[j] = temp; + } + } + + const newSlotDistribution = shuffledProxies.map((proxy, index) => { + const slotLength = Math.floor(16384 / shuffledProxies.length); + const from = index * slotLength; + const to = index === shuffledProxies.length - 1 ? 16383 : from + slotLength - 1; + return { proxy, slotRanges: `${from}-${to}` }; + }); + + const customShuffledInterceptor = { + name: "cluster-simulation-interceptor", + fn: async (data: Buffer, next: (data: Buffer) => Promise, state: { invokeCount: number; matchCount: number }) => { + state.invokeCount++; + if (data.toString().toLowerCase() !== "*2\r\n$7\r\ncluster\r\n$5\r\nslots\r\n") { + return next(data); + } + state.matchCount++; + const mapping = newSlotDistribution.map(({ proxy, slotRanges }) => { + const [from, to] = slotRanges.split("-").map(Number); + const id = `proxy-id-${proxy.config.listenPort}`; + return `*3\r\n:${from}\r\n:${to}\r\n*3\r\n$${proxy.config.listenHost.length}\r\n${proxy.config.listenHost}\r\n:${proxy.config.listenPort}\r\n$${id.length}\r\n${id}\r\n`; + }); + return Buffer.from(`*${newSlotDistribution.length}\r\n${mapping.join("")}`); + }, + }; + + for (const proxy of proxyStore.proxies) { + proxy.addGlobalInterceptor(customShuffledInterceptor); + } + + for (const { proxy, slotRanges } of oldSlotDistribution) { + const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); + proxy.sendToAllClients(sMigratingBuffer); + } + + await delay(5000); + + for (const { proxy: sourceProxy } of oldSlotDistribution) { + const migratedSlots = newSlotDistribution.map(({ proxy, slotRanges }) => ({ + targetNode: { host: proxy.config.listenHost, port: proxy.config.listenPort }, + slotRanges, + })); + const sMigratedBuffer = buildSMigratedNotification(migratedSlots, getNextSequenceId()); + sourceProxy.sendToAllClients(sMigratedBuffer); + } +} diff --git a/src/app.ts b/src/app.ts index 0efcbf0..2de1c94 100644 --- a/src/app.ts +++ b/src/app.ts @@ -10,10 +10,15 @@ import { RedisProxy, type SendResult, } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy.ts"; +import { executeAction } from "./actions/index.ts"; import applyDefaultInterceptors from "./default_interceptors/index.ts"; import ProxyStore, { makeId } from "./proxy-store.ts"; -import applyPredefinedScenario from "./scenarios/index.ts"; import { + type ActionRecord, + type ActionTrigger, + type ListActionTriggersResponse, + actionIdParamSchema, + actionRequestSchema, connectionIdsQuerySchema, type ExtendedProxyConfig, encodingSchema, @@ -21,9 +26,9 @@ import { interceptorSchema, paramSchema, parseBuffer, - predefinedScenarioParamSchema, proxyConfigSchema, - scenarioSchema, + slotMigrateEffectSchema, + type SlotMigrateEffect, } from "./util.ts"; const startNewProxy = (config: ProxyConfig) => { @@ -148,42 +153,6 @@ export function createApp(testConfig?: ExtendedProxyConfig) { return c.json({ success, connectionId }); }); - app.post("/scenarios", zValidator("json", scenarioSchema), async (c) => { - const { responses, encoding } = c.req.valid("json"); - - const responsesBuffers = responses.map((response) => parseBuffer(response, encoding)); - let currentIndex = 0; - - const scenarioInterceptor: InterceptorDescription = { - name: "scenario-interceptor", - fn: async (data: Buffer, next: Next, state: InterceptorState): Promise => { - state.invokeCount++; - if (currentIndex < responsesBuffers.length) { - state.matchCount++; - const response = responsesBuffers[currentIndex] as Buffer; - currentIndex++; - return response; - } - return await next(data); - }, - }; - - for (const proxy of proxyStore.proxies) { - proxy.addGlobalInterceptor(scenarioInterceptor); - } - - return c.json({ success: true, totalResponses: responses.length }); - }); - - app.post( - "/scenarios/predefined/:scenario", - zValidator("param", predefinedScenarioParamSchema), - async (c) => { - const { scenario } = c.req.valid("param"); - return await applyPredefinedScenario(scenario, c, proxyStore, config); - }, - ); - app.post("/interceptors", zValidator("json", interceptorSchema), async (c) => { const { name, match, response, encoding } = c.req.valid("json"); @@ -209,5 +178,142 @@ export function createApp(testConfig?: ExtendedProxyConfig) { return c.json({ success: true, name }); }); + // In-memory action storage + const actionStore = new Map(); + + // Generate unique action ID + const generateActionId = (): string => { + return `action-${Date.now()}-${Math.random().toString(36).substring(2, 9)}`; + }; + + // POST /action - Submit an action + app.post("/action", zValidator("json", actionRequestSchema), async (c) => { + const { type, parameters } = c.req.valid("json"); + + const actionId = generateActionId(); + const actionRecord: ActionRecord = { + id: actionId, + type, + parameters, + status: "pending", + submittedAt: new Date(), + error: null, + output: null, + }; + + actionStore.set(actionId, actionRecord); + + // Execute the action asynchronously + executeAction(type, parameters, proxyStore, config) + .then((result) => { + actionRecord.status = result.status; + actionRecord.output = "Done"; + actionRecord.error = result.error; + }) + .catch((error) => { + actionRecord.status = "failed"; + actionRecord.error = error instanceof Error ? error.message : String(error); + }); + + return c.json({ action_id: actionId }); + }); + + // GET /action/:action_id - Get action status + app.get("/action/:action_id", zValidator("param", actionIdParamSchema), (c) => { + const { action_id } = c.req.valid("param"); + + const action = actionStore.get(action_id); + if (!action) { + return c.json({ error: "Action not found" }, 404); + } + + return c.json({ + status: action.status, + error: action.error, + output: action.output, + }); + }); + + // Hardcoded triggers for each effect + const triggersMap: Record = { + add: [ + { + name: "add-node-trigger-1", + description: "Trigger when a new node is added to the cluster", + requirements: [ + { dbconfig: {}, cluster: { minNodes: 3 }, description: "Requires at least 3 shards and 3 nodes" }, + ], + }, + { + name: "add-node-trigger-2", + description: "Trigger for rebalancing after node addition", + requirements: [ + { dbconfig: {}, cluster: { healthy: true }, description: "Requires replication enabled and healthy cluster" }, + ], + }, + ], + remove: [ + { + name: "remove-node-trigger-1", + description: "Trigger when a node is removed from the cluster", + requirements: [ + { dbconfig: {}, cluster: { minNodes: 2 }, description: "Requires at least 2 shards and 2 nodes" }, + ], + }, + { + name: "remove-node-trigger-2", + description: "Trigger for slot migration before node removal", + requirements: [ + { dbconfig: {}, cluster: { noFailover: true }, description: "Requires persistence and no ongoing failover" }, + ], + }, + ], + "remove-add": [ + { + name: "remove-add-trigger-1", + description: "Trigger for combined remove and add operation", + requirements: [ + { dbconfig: {}, cluster: { minNodes: 3 }, description: "Requires at least 3 shards and 3 nodes" }, + ], + }, + { + name: "remove-add-trigger-2", + description: "Trigger for atomic node replacement", + requirements: [ + { dbconfig: {}, cluster: { quorum: true }, description: "Requires replication and quorum" }, + ], + }, + ], + "slot-shuffle": [ + { + name: "slot-shuffle-trigger-1", + description: "Trigger for redistributing slots across nodes", + requirements: [ + { dbconfig: {}, cluster: { balanced: false }, description: "Requires at least 2 shards and unbalanced cluster" }, + ], + }, + { + name: "slot-shuffle-trigger-2", + description: "Trigger for optimizing slot distribution", + requirements: [ + { dbconfig: {}, cluster: { healthy: true }, description: "Requires auto-balance enabled and healthy cluster" }, + ], + }, + ], + }; + + // GET /slot-migrate - List action triggers for an effect + app.get("/slot-migrate", zValidator("query", slotMigrateEffectSchema), (c) => { + const { effect } = c.req.valid("query"); + + const response: ListActionTriggersResponse = { + effect, + cluster: { index: 0, nodes: proxyStore.nodeIds.length }, + triggers: triggersMap[effect], + }; + + return c.json(response); + }); + return { app, proxy: proxyStore.proxies[0] as RedisProxy, config }; } diff --git a/src/scenarios/slot-shuffle.ts b/src/scenarios/slot-shuffle.ts index 60d7189..ab06138 100644 --- a/src/scenarios/slot-shuffle.ts +++ b/src/scenarios/slot-shuffle.ts @@ -6,7 +6,6 @@ import type { ExtendedProxyConfig } from "../util"; import { buildSMigratedNotification, buildSMigratingNotification, - createCustomClusterSlotsInterceptor, getSlotRangesForProxy, } from "./helpers"; import { getNextSequenceId } from "./sequence-gen"; diff --git a/src/util.ts b/src/util.ts index 87e5a24..203f45e 100644 --- a/src/util.ts +++ b/src/util.ts @@ -24,11 +24,6 @@ export const connectionIdsQuerySchema = z.object({ encoding: z.enum(["base64", "raw"]).default("base64"), }); -export const scenarioSchema = z.object({ - responses: z.array(z.string()).min(1, "At least one response is required"), - encoding: z.enum(["base64", "raw"]).default("base64"), -}); - export const interceptorSchema = z.object({ name: z.string(), encoding: z.enum(["raw", "base64"]), @@ -40,6 +35,80 @@ export const predefinedScenarioParamSchema = z.object({ scenario: z.enum(["remove-add", "remove", "add", "slot-shuffle"]), }); +export const slotMigrateEffectSchema = z.object({ + effect: z.enum(["add", "remove", "remove-add", "slot-shuffle"]), +}); + +export type SlotMigrateEffect = z.infer["effect"]; + +export interface ActionTriggerRequirement { + dbconfig: unknown; + cluster: unknown; + description: string; +} + +export interface ActionTrigger { + name: string; + description: string; + requirements: ActionTriggerRequirement[]; +} + +export interface ListActionTriggersResponse { + effect: string; + cluster: { index: number; nodes: number }; + triggers: ActionTrigger[]; +} + +// Action types matching re_fault_injector ActionType enum +export const actionTypeSchema = z.enum([ + "dmc_restart", + "failover", + "reshard", + "sequence_of_actions", + "network_failure", + "execute_rlutil_command", + "execute_rladmin_command", + "enable_entraid", + "upgrade", + "wait", + "migrate", + "bind", + "update_cluster_config", + "delete_database", + "create_database", + "shard_failure", + "node_failure", + "node_remove", + "proxy_failure", + "cluster_failure", + "slot_migrate", +]); + +export type ActionType = z.infer; + +export const actionRequestSchema = z.object({ + type: actionTypeSchema, + parameters: z.record(z.string(), z.unknown()), +}); + +export type ActionRequest = z.infer; + +export const actionIdParamSchema = z.object({ + action_id: z.string(), +}); + +export type ActionStatus = "pending" | "running" | "success" | "failed" | "unknown"; + +export interface ActionRecord { + id: string; + type: ActionType; + parameters: Record; + status: ActionStatus; + submittedAt: Date; + error?: string | null; + output?: unknown; +} + export function parseBuffer(data: string, encoding: "base64" | "raw"): Buffer { switch (encoding) { case "base64": From a4e82bac774e208694eadf982a921f90e87bd1eb Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Thu, 10 Sep 2026 14:13:46 +0300 Subject: [PATCH 4/5] feat(actions): mimic the Fault Injector API Make the proxy's action API a drop-in stand-in for the Fault Injector (re_fault_injector) so FI clients work against it unchanged: - Implement reset_cluster: restore the initial proxy topology and reapply default interceptors. - Implement create_database: size the proxy cluster to database_config.shards_count and return raw_endpoints, username, password, tls and bdb_id in the action output, matching what FI clients parse. - Add GET /action to list submitted actions. - Report "running" while an action executes and return structured action output instead of the hardcoded "Done". - Generate FI-shaped triggers for GET /slot-migrate: real trigger names (migrate, maintenance_mode, failover), descriptions, and full dbconfig requirements with ext-ip/ext-hostname names, mirroring the FI's TRIGGER_DEFINITIONS and dbconfig templates. - Broadcast SMIGRATING/SMIGRATED to clients on every node during remove and remove-add, so each cluster connection observes the migration. - Push SMIGRATING to connections opened during an active migration, matching how a real cluster treats new connections mid-migration. - Complete the action type enum to FI parity (reset_cluster, topology_change_standalone, network_latency, wait_for_database_active, collect_debuginfo). - Make migration delays env-overridable (MIGRATION_DELAY_MS, COMPLETION_DELAY_MS) so tests run fast. - Remove the dead /scenarios modules and their stale tests; cover the action API with new tests instead. Co-Authored-By: Claude Fable 5 --- src/actions.test.ts | 271 ++++++++++++++++++++++++++++++++++ src/actions/index.ts | 233 +++++++++++++++++++++++++---- src/actions/triggers.ts | 118 +++++++++++++++ src/app.ts | 89 ++--------- src/scenarios.test.ts | 153 ------------------- src/scenarios/add.ts | 119 --------------- src/scenarios/index.ts | 37 ----- src/scenarios/remove-add.ts | 121 --------------- src/scenarios/remove.ts | 125 ---------------- src/scenarios/slot-shuffle.ts | 150 ------------------- src/util.ts | 14 +- 11 files changed, 618 insertions(+), 812 deletions(-) create mode 100644 src/actions.test.ts create mode 100644 src/actions/triggers.ts delete mode 100644 src/scenarios.test.ts delete mode 100644 src/scenarios/add.ts delete mode 100644 src/scenarios/index.ts delete mode 100644 src/scenarios/remove-add.ts delete mode 100644 src/scenarios/remove.ts delete mode 100644 src/scenarios/slot-shuffle.ts diff --git a/src/actions.test.ts b/src/actions.test.ts new file mode 100644 index 0000000..e760a2b --- /dev/null +++ b/src/actions.test.ts @@ -0,0 +1,271 @@ +import { afterAll, beforeAll, describe, expect, test } from "bun:test"; +import type { Socket } from "bun"; + +import { getFreePortNumber } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy.ts"; +import { createApp } from "./app"; +import createMockRedisServer from "./mock-server"; + +// Shrink the migration windows so the suite stays fast. The action code reads +// these at execution time. +process.env.MIGRATION_DELAY_MS = "300"; +process.env.COMPLETION_DELAY_MS = "100"; + +interface ActionStatusResponse { + status: string; + error: unknown; + output: unknown; +} + +describe("Fault Injector action API", () => { + let app: any; + let mockRedisServer: ReturnType; + let listenPort: number; + let targetPort: number; + + const postAction = async (body: unknown): Promise => { + return app.request("/action", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify(body), + }); + }; + + const submitAction = async (body: unknown): Promise => { + const res = await postAction(body); + expect(res.status).toBe(200); + const { action_id } = (await res.json()) as { action_id: string }; + expect(action_id).toBeString(); + return action_id; + }; + + const waitForAction = async (actionId: string): Promise => { + const deadline = Date.now() + 10_000; + while (Date.now() < deadline) { + const res = await app.request(`/action/${actionId}`); + expect(res.status).toBe(200); + const action = (await res.json()) as ActionStatusResponse; + if (action.status !== "pending" && action.status !== "running") { + return action; + } + await new Promise((resolve) => setTimeout(resolve, 25)); + } + throw new Error(`Timeout waiting for action ${actionId}`); + }; + + const runAction = async (body: unknown): Promise => { + return waitForAction(await submitAction(body)); + }; + + const getNodeIds = async (): Promise => { + const res = await app.request("/nodes"); + return (await res.json()).ids; + }; + + const resetCluster = async () => { + const result = await runAction({ type: "reset_cluster", parameters: {} }); + expect(result.status).toBe("success"); + }; + + beforeAll(async () => { + listenPort = await getFreePortNumber(); + targetPort = await getFreePortNumber(); + + mockRedisServer = createMockRedisServer(targetPort); + + const appInstance = createApp({ + listenPort: [listenPort], + listenHost: "127.0.0.1", + targetHost: "127.0.0.1", + targetPort: targetPort, + timeout: 30000, + enableLogging: false, + apiPort: 3002, + }); + app = appInstance.app; + + await new Promise((resolve) => setTimeout(resolve, 200)); + }); + + afterAll(async () => { + const res = await app.request("/nodes"); + const { ids } = await res.json(); + for (const id of ids) { + await app.request(`/nodes/${encodeURIComponent(id)}`, { method: "DELETE" }); + } + mockRedisServer?.stop(true); + }); + + test("POST /action rejects unknown action types", async () => { + const res = await postAction({ type: "bogus", parameters: {} }); + expect(res.status).toBe(400); + }); + + test("GET /slot-migrate returns FI-shaped triggers", async () => { + const res = await app.request("/slot-migrate?effect=remove"); + expect(res.status).toBe(200); + + const body = await res.json(); + expect(body.effect).toBe("remove"); + expect(body.cluster.index).toBe(0); + expect(body.cluster.nodes).toBeNumber(); + + const names = body.triggers.map((trigger: { name: string }) => trigger.name); + expect(names).toEqual(["migrate", "maintenance_mode", "failover"]); + + for (const trigger of body.triggers) { + expect(trigger.description).toBeString(); + expect(trigger.requirements.length).toBeGreaterThan(0); + for (const requirement of trigger.requirements) { + expect(requirement.dbconfig.name).toBeString(); + expect(requirement.dbconfig.name).toContain("sm-remove-"); + expect(requirement.dbconfig.shards_count).toBeNumber(); + expect(requirement.cluster.min_nodes).toBe(3); + expect(requirement.description).toBeString(); + } + } + + // The --db=ext-hostname CLI filter of the scenario tests matches on this + const dbNames = body.triggers.flatMap( + (trigger: { requirements: { dbconfig: { name: string } }[] }) => + trigger.requirements.map((requirement) => requirement.dbconfig.name), + ); + expect(dbNames.some((name: string) => name.includes("ext-ip"))).toBe(true); + expect(dbNames.some((name: string) => name.includes("ext-hostname"))).toBe(true); + }); + + test("GET /slot-migrate requires an effect", async () => { + const res = await app.request("/slot-migrate"); + expect(res.status).toBe(400); + }); + + test("create_database sizes the cluster and returns connection info", async () => { + const result = await runAction({ + type: "create_database", + parameters: { + cluster_index: 0, + database_config: { name: "sm-remove-migrate-ext-ip", shards_count: 3 }, + }, + }); + + expect(result.status).toBe("success"); + const output = result.output as { + bdb_id: number; + username: string; + password: string; + tls: boolean; + raw_endpoints: { dns_name: string; port: number }[]; + }; + expect(output.bdb_id).toBeNumber(); + expect(output.tls).toBe(false); + expect(output.raw_endpoints.length).toBe(3); + expect(output.raw_endpoints[0]?.dns_name).toBe("127.0.0.1"); + expect(output.raw_endpoints[0]?.port).toBeNumber(); + + expect((await getNodeIds()).length).toBe(3); + }); + + test("reset_cluster restores the initial topology", async () => { + await resetCluster(); + expect((await getNodeIds()).length).toBe(1); + }); + + test("GET /action lists submitted actions", async () => { + const actionId = await submitAction({ type: "wait", parameters: {} }); + await waitForAction(actionId); + + const res = await app.request("/action"); + expect(res.status).toBe(200); + const { actions } = await res.json(); + const entry = actions.find((action: { job_id: string }) => action.job_id === actionId); + expect(entry).toBeDefined(); + expect(entry.action_type).toBe("wait"); + expect(entry.submitted_at).toBeString(); + }); + + test("slot_migrate rejects an invalid effect", async () => { + const result = await runAction({ + type: "slot_migrate", + parameters: { effect: "explode", cluster_index: 0 }, + }); + expect(result.status).toBe("failed"); + expect(String(result.error)).toContain("Invalid effect"); + }); + + test("slot_migrate remove notifies all clients and new connections", async () => { + await resetCluster(); + const createResult = await runAction({ + type: "create_database", + parameters: { + cluster_index: 0, + database_config: { name: "sm-remove-migrate-ext-ip", shards_count: 3 }, + }, + }); + const { raw_endpoints } = createResult.output as { + raw_endpoints: { dns_name: string; port: number }[]; + }; + expect(raw_endpoints.length).toBe(3); + + const connectAndCollect = async (port: number) => { + const chunks: Buffer[] = []; + let socket: Socket | undefined; + await new Promise((resolve, reject) => { + Bun.connect({ + hostname: "127.0.0.1", + port, + socket: { + open(openedSocket) { + socket = openedSocket; + resolve(); + }, + data(_socket, data) { + chunks.push(Buffer.from(data)); + }, + error(_socket, error) { + reject(error); + }, + close() {}, + }, + }); + }); + return { + received: () => Buffer.concat(chunks).toString(), + close: () => socket?.end(), + }; + }; + + const clients = await Promise.all( + raw_endpoints.map((endpoint) => connectAndCollect(endpoint.port)), + ); + // Let the proxies register the connections + await new Promise((resolve) => setTimeout(resolve, 100)); + + const actionId = await submitAction({ + type: "slot_migrate", + parameters: { effect: "remove", cluster_index: 0, trigger: "migrate", bdb_id: "1" }, + }); + + // SMIGRATING is broadcast right away; the migration window is 300ms + await new Promise((resolve) => setTimeout(resolve, 100)); + for (const client of clients) { + expect(client.received()).toContain("SMIGRATING"); + } + + // A connection opened during the migration window gets SMIGRATING too + const lateClient = await connectAndCollect(raw_endpoints[0]?.port as number); + await new Promise((resolve) => setTimeout(resolve, 150)); + expect(lateClient.received()).toContain("SMIGRATING"); + + const result = await waitForAction(actionId); + expect(result.status).toBe("success"); + + // One node was removed and the survivors got SMIGRATED + expect((await getNodeIds()).length).toBe(2); + const migrated = clients.filter((client) => client.received().includes("SMIGRATED")); + expect(migrated.length).toBeGreaterThan(0); + + lateClient.close(); + for (const client of clients) { + client.close(); + } + }); +}); diff --git a/src/actions/index.ts b/src/actions/index.ts index 9b46467..7a61485 100644 --- a/src/actions/index.ts +++ b/src/actions/index.ts @@ -1,7 +1,10 @@ -import type { ProxyConfig } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy"; +import { + type ProxyConfig, + RedisProxy, +} from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy"; +import applyDefaultInterceptors from "../default_interceptors/index"; import type ProxyStore from "../proxy-store"; import { makeId } from "../proxy-store"; -import type { ActionType, ExtendedProxyConfig } from "../util"; import { addNode, buildSMigratedNotification, @@ -10,8 +13,10 @@ import { findNextAvailablePort, getSlotRangesForProxy, pickRandom, + sendToAllClients, } from "../scenarios/helpers"; import { getNextSequenceId } from "../scenarios/sequence-gen"; +import type { ActionType, ExtendedProxyConfig } from "../util"; // Effect types matching Python MigrateEffect enum export type SlotMigrateEffect = "remove-add" | "remove" | "add" | "slot-shuffle"; @@ -26,10 +31,71 @@ export interface SlotMigrateParams { export interface ActionExecutionResult { status: "success" | "failed"; error?: string | null; + output?: unknown; } const delay = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); +// The window between SMIGRATING and SMIGRATED, and the settle time before a +// removed node is stopped. Overridable so tests do not wait several seconds. +const migrationDelayMs = () => Number(process.env.MIGRATION_DELAY_MS ?? 5000); +const completionDelayMs = () => Number(process.env.COMPLETION_DELAY_MS ?? 2000); + +// How long to wait after a client connects before pushing SMIGRATING to it, +// so the push lands after the client finished its HELLO handshake. +const NEW_CONNECTION_PUSH_DELAY_MS = 25; + +/** + * While a migration is active, push the SMIGRATING notification to every + * client that connects, matching how a real cluster treats connections + * opened while a migration is in progress. + * Returns a function that stops the notifications. + */ +function pushToNewConnections(proxyStore: ProxyStore, buffer: Buffer): () => void { + const subscriptions = proxyStore.proxies.map((proxy) => { + const listener = (connection: { id: string }) => { + setTimeout(() => proxy.sendToClient(connection.id, buffer), NEW_CONNECTION_PUSH_DELAY_MS); + }; + proxy.on("connection", listener); + return () => proxy.off("connection", listener); + }); + return () => { + for (const unsubscribe of subscriptions) unsubscribe(); + }; +} + +async function startNode(proxyStore: ProxyStore, config: ProxyConfig): Promise { + const proxy = new RedisProxy(config); + // Without a listener, the 'error' the proxy emits alongside a failed + // start() escapes the EventEmitter and crashes the process. + proxy.on("error", (error: Error) => console.error("[proxy]", error.message)); + await proxy.start(); + proxyStore.add(makeId(config.targetHost, config.targetPort, config.listenPort), proxy); + return proxy; +} + +/** Next free listen port, never colliding with the backend target port. */ +function nextListenPort(proxyStore: ProxyStore, config: ExtendedProxyConfig): number { + let port = + proxyStore.proxies.length > 0 + ? findNextAvailablePort(proxyStore.proxies) + : Math.max(...config.listenPort); + while (port === config.targetPort) port++; + return port; +} + +async function removeNode(proxyStore: ProxyStore, proxy: RedisProxy): Promise { + const { targetHost, targetPort, listenPort } = proxy.config; + await proxyStore.delete(makeId(targetHost, targetPort, listenPort)); +} + +function refreshClusterSlots(proxyStore: ProxyStore): void { + const interceptor = createCustomClusterSlotsInterceptor(proxyStore.proxies); + for (const proxy of proxyStore.proxies) { + proxy.addGlobalInterceptor(interceptor); + } +} + /** * Execute an action based on its type and parameters */ @@ -42,11 +108,85 @@ export async function executeAction( switch (actionType) { case "slot_migrate": return executeSlotMigrate(parameters as unknown as SlotMigrateParams, proxyStore, config); + case "reset_cluster": + return executeResetCluster(proxyStore, config); + case "create_database": + return executeCreateDatabase(parameters, proxyStore, config); default: return { status: "success" }; } } +/** + * Restore the initial proxy topology and drop interceptors added by earlier + * actions. Test harnesses call this before every test. + */ +async function executeResetCluster( + proxyStore: ProxyStore, + config: ExtendedProxyConfig, +): Promise { + try { + for (const id of proxyStore.nodeIds) { + await proxyStore.delete(id); + } + for (const port of config.listenPort) { + await startNode(proxyStore, { ...config, listenPort: port }); + } + if (config.defaultInterceptors) { + applyDefaultInterceptors(config.defaultInterceptors, proxyStore); + } + return { status: "success" }; + } catch (error) { + return { status: "failed", error: error instanceof Error ? error.message : String(error) }; + } +} + +/** + * Size the proxy cluster to the requested shards_count and return connection + * info in the shape Fault Injector clients expect: they read + * raw_endpoints[0], username, password, tls and bdb_id from the output. + */ +async function executeCreateDatabase( + parameters: Record, + proxyStore: ProxyStore, + config: ExtendedProxyConfig, +): Promise { + try { + const databaseConfig = (parameters.database_config ?? {}) as Record; + const shardsCount = + typeof databaseConfig.shards_count === "number" && databaseConfig.shards_count > 0 + ? databaseConfig.shards_count + : Math.max(proxyStore.proxies.length, 1); + + while (proxyStore.proxies.length > shardsCount) { + const proxy = proxyStore.proxies.at(-1); + if (!proxy) break; + await removeNode(proxyStore, proxy); + } + while (proxyStore.proxies.length < shardsCount) { + await startNode(proxyStore, { ...config, listenPort: nextListenPort(proxyStore, config) }); + } + + refreshClusterSlots(proxyStore); + + return { + status: "success", + output: { + bdb_id: 1, + username: "", + password: "", + tls: false, + raw_endpoints: proxyStore.proxies.map((proxy) => ({ + dns_name: proxy.config.listenHost, + port: proxy.config.listenPort, + })), + }, + }; + } catch (error) { + return { status: "failed", error: error instanceof Error ? error.message : String(error) }; + } +} + /** * Execute slot_migrate action - handles all scenario effects */ @@ -63,7 +203,10 @@ async function executeSlotMigrate( const validEffects: SlotMigrateEffect[] = ["remove-add", "remove", "add", "slot-shuffle"]; if (!validEffects.includes(effect)) { - return { status: "failed", error: `Invalid effect: ${effect}. Must be one of: ${validEffects.join(", ")}` }; + return { + status: "failed", + error: `Invalid effect: ${effect}. Must be one of: ${validEffects.join(", ")}`, + }; } try { @@ -87,7 +230,10 @@ async function executeSlotMigrate( } } -async function executeRemoveAddEffect(proxyStore: ProxyStore, config: ExtendedProxyConfig): Promise { +async function executeRemoveAddEffect( + proxyStore: ProxyStore, + config: ExtendedProxyConfig, +): Promise { const allProxies = proxyStore.proxies; if (allProxies.length === 0) { throw new Error("No proxies available to select from"); @@ -99,7 +245,7 @@ async function executeRemoveAddEffect(proxyStore: ProxyStore, config: ExtendedPr } const slotRanges = getSlotRangesForProxy(proxyToBeRemoved, allProxies); - const newPort = findNextAvailablePort(allProxies); + const newPort = nextListenPort(proxyStore, config); const newProxyConfig: ProxyConfig = { ...config, listenPort: newPort }; const { proxy: newProxy } = addNode(proxyStore, newProxyConfig); @@ -111,23 +257,32 @@ async function executeRemoveAddEffect(proxyStore: ProxyStore, config: ExtendedPr } const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); - proxyToBeRemoved.sendToAllClients(sMigratingBuffer); + sendToAllClients(proxyStore, sMigratingBuffer); + const stopNotifying = pushToNewConnections(proxyStore, sMigratingBuffer); - await delay(5000); + await delay(migrationDelayMs()); const sMigratedBuffer = buildSMigratedNotification( - [{ targetNode: { host: newProxy.config.listenHost, port: newProxy.config.listenPort }, slotRanges }], + [ + { + targetNode: { host: newProxy.config.listenHost, port: newProxy.config.listenPort }, + slotRanges, + }, + ], getNextSequenceId(), ); - proxyToBeRemoved.sendToAllClients(sMigratedBuffer); + sendToAllClients(proxyStore, sMigratedBuffer); + stopNotifying(); - await delay(2000); + await delay(completionDelayMs()); - const { targetHost, targetPort, listenPort } = proxyToBeRemoved.config; - await proxyStore.delete(makeId(targetHost, targetPort, listenPort)); + await removeNode(proxyStore, proxyToBeRemoved); } -async function executeRemoveEffect(proxyStore: ProxyStore, _config: ExtendedProxyConfig): Promise { +async function executeRemoveEffect( + proxyStore: ProxyStore, + _config: ExtendedProxyConfig, +): Promise { const allProxies = proxyStore.proxies; if (allProxies.length === 0) { throw new Error("No proxies available to select from"); @@ -155,24 +310,28 @@ async function executeRemoveEffect(proxyStore: ProxyStore, _config: ExtendedProx } const sMigratingBuffer = buildSMigratingNotification(removedNodeSlotRanges, getNextSequenceId()); - proxyToBeRemoved.sendToAllClients(sMigratingBuffer); + sendToAllClients(proxyStore, sMigratingBuffer); + const stopNotifying = pushToNewConnections(proxyStore, sMigratingBuffer); - await delay(5000); + await delay(migrationDelayMs()); const migratedSlots = newSlotDistribution.map(({ proxy, slotRanges }) => ({ targetNode: { host: proxy.config.listenHost, port: proxy.config.listenPort }, slotRanges, })); const sMigratedBuffer = buildSMigratedNotification(migratedSlots, getNextSequenceId()); - proxyToBeRemoved.sendToAllClients(sMigratedBuffer); + sendToAllClients(proxyStore, sMigratedBuffer); + stopNotifying(); - await delay(2000); + await delay(completionDelayMs()); - const { targetHost, targetPort, listenPort } = proxyToBeRemoved.config; - await proxyStore.delete(makeId(targetHost, targetPort, listenPort)); + await removeNode(proxyStore, proxyToBeRemoved); } -async function executeAddEffect(proxyStore: ProxyStore, config: ExtendedProxyConfig): Promise { +async function executeAddEffect( + proxyStore: ProxyStore, + config: ExtendedProxyConfig, +): Promise { const allProxies = proxyStore.proxies; if (allProxies.length === 0) { throw new Error("No proxies available"); @@ -183,7 +342,7 @@ async function executeAddEffect(proxyStore: ProxyStore, config: ExtendedProxyCon slotRanges: getSlotRangesForProxy(proxy, allProxies), })); - const newPort = findNextAvailablePort(allProxies); + const newPort = nextListenPort(proxyStore, config); const newProxyConfig: ProxyConfig = { ...config, listenPort: newPort }; const { proxy: newProxy } = addNode(proxyStore, newProxyConfig); @@ -194,24 +353,35 @@ async function executeAddEffect(proxyStore: ProxyStore, config: ExtendedProxyCon proxy.addGlobalInterceptor(clusterSlotsInterceptor); } + let sMigratingBuffer: Buffer = Buffer.alloc(0); for (const { proxy, slotRanges } of oldSlotDistribution) { - const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); + sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); proxy.sendToAllClients(sMigratingBuffer); } + const stopNotifying = pushToNewConnections(proxyStore, sMigratingBuffer); - await delay(5000); + await delay(migrationDelayMs()); const newNodeSlotRanges = getSlotRangesForProxy(newProxy, allProxiesWithNew); for (const { proxy } of oldSlotDistribution) { const sMigratedBuffer = buildSMigratedNotification( - [{ targetNode: { host: newProxy.config.listenHost, port: newProxy.config.listenPort }, slotRanges: newNodeSlotRanges }], + [ + { + targetNode: { host: newProxy.config.listenHost, port: newProxy.config.listenPort }, + slotRanges: newNodeSlotRanges, + }, + ], getNextSequenceId(), ); proxy.sendToAllClients(sMigratedBuffer); } + stopNotifying(); } -async function executeSlotShuffleEffect(proxyStore: ProxyStore, _config: ExtendedProxyConfig): Promise { +async function executeSlotShuffleEffect( + proxyStore: ProxyStore, + _config: ExtendedProxyConfig, +): Promise { const allProxies = proxyStore.proxies; if (allProxies.length === 0) { throw new Error("No proxies available"); @@ -246,7 +416,11 @@ async function executeSlotShuffleEffect(proxyStore: ProxyStore, _config: Extende const customShuffledInterceptor = { name: "cluster-simulation-interceptor", - fn: async (data: Buffer, next: (data: Buffer) => Promise, state: { invokeCount: number; matchCount: number }) => { + fn: async ( + data: Buffer, + next: (data: Buffer) => Promise, + state: { invokeCount: number; matchCount: number }, + ) => { state.invokeCount++; if (data.toString().toLowerCase() !== "*2\r\n$7\r\ncluster\r\n$5\r\nslots\r\n") { return next(data); @@ -265,12 +439,14 @@ async function executeSlotShuffleEffect(proxyStore: ProxyStore, _config: Extende proxy.addGlobalInterceptor(customShuffledInterceptor); } + let sMigratingBuffer: Buffer = Buffer.alloc(0); for (const { proxy, slotRanges } of oldSlotDistribution) { - const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); + sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); proxy.sendToAllClients(sMigratingBuffer); } + const stopNotifying = pushToNewConnections(proxyStore, sMigratingBuffer); - await delay(5000); + await delay(migrationDelayMs()); for (const { proxy: sourceProxy } of oldSlotDistribution) { const migratedSlots = newSlotDistribution.map(({ proxy, slotRanges }) => ({ @@ -280,4 +456,5 @@ async function executeSlotShuffleEffect(proxyStore: ProxyStore, _config: Extende const sMigratedBuffer = buildSMigratedNotification(migratedSlots, getNextSequenceId()); sourceProxy.sendToAllClients(sMigratedBuffer); } + stopNotifying(); } diff --git a/src/actions/triggers.ts b/src/actions/triggers.ts new file mode 100644 index 0000000..eba465f --- /dev/null +++ b/src/actions/triggers.ts @@ -0,0 +1,118 @@ +import type { ActionTrigger, ActionTriggerRequirement, SlotMigrateEffect } from "../util"; + +// Mirrors re_fault_injector TRIGGER_DEFINITIONS so clients can select +// triggers by their real names (migrate, maintenance_mode, failover). +const TRIGGER_DEFINITIONS: Record = { + "remove-add": [ + { + name: "migrate", + description: "Use rladmin migrate to move all shards from source node to empty node", + }, + { + name: "maintenance_mode", + description: "Put source node in maintenance mode, shards auto-migrate to other nodes", + }, + { + name: "failover", + description: "Trigger failover to swap master/replica roles (requires replication)", + }, + ], + remove: [ + { + name: "migrate", + description: "Use rladmin migrate to move all shards from source node to existing node", + }, + { + name: "maintenance_mode", + description: "Put source node in maintenance mode, shards auto-migrate to other nodes", + }, + { + name: "failover", + description: "Trigger failover to swap master/replica roles (requires replication)", + }, + ], + add: [ + { name: "migrate", description: "Use rladmin migrate to move one shard to empty node" }, + { + name: "failover", + description: "Trigger failover to swap master/replica roles (requires replication)", + }, + ], + "slot-shuffle": [ + { + name: "migrate", + description: "Use rladmin migrate to move one shard between existing nodes", + }, + { + name: "failover", + description: "Trigger failover to swap master/replica roles (requires replication)", + }, + ], +}; + +// The two default (external) OSS Cluster API combinations the FI returns. +const IP_TYPES = [ + { ipType: "external", endpointType: "ip", suffix: "ext-ip" }, + { ipType: "external", endpointType: "hostname", suffix: "ext-hostname" }, +]; + +const MIN_NODES = 3; +const BASE_PORT = 13000; + +// Mirrors re_fault_injector _calculate_shards_count. For the proxy, +// shards_count doubles as the number of proxy nodes create_database spins up. +function calculateShardsCount(effect: SlotMigrateEffect, trigger: string, nodeCount: number) { + switch (effect) { + case "remove-add": + return trigger === "failover" ? 1 : nodeCount - 1; + case "remove": + return nodeCount; + case "add": + return trigger === "failover" ? 2 : nodeCount; + case "slot-shuffle": + return nodeCount * 2; + } +} + +function calculatePlacement(effect: SlotMigrateEffect, trigger: string) { + if (effect === "add") return "dense"; + if (effect === "remove-add" && trigger === "maintenance_mode") return "dense"; + return "sparse"; +} + +export function generateTriggersForEffect( + effect: SlotMigrateEffect, + nodeCount: number, +): ActionTrigger[] { + // The FI derives shards_count from the RE cluster size and needs >= 3 nodes + // for quorum. The proxy has no quorum and can add nodes freely, so size the + // dbconfigs as if the cluster had at least MIN_NODES. + const effectiveNodes = Math.max(nodeCount, MIN_NODES); + + return TRIGGER_DEFINITIONS[effect].map((trigger) => ({ + name: trigger.name, + description: trigger.description, + requirements: IP_TYPES.map(({ ipType, endpointType, suffix }, index) => { + const requirement: ActionTriggerRequirement = { + dbconfig: { + name: `sm-${effect}-${trigger.name.replace(/_/g, "-")}-${suffix}`, + port: BASE_PORT + index, + memory_size: 134217728, + eviction_policy: "volatile-lru", + sharding: true, + oss_cluster: true, + proxy_policy: "all-master-shards", + shards_count: calculateShardsCount(effect, trigger.name, effectiveNodes), + shards_placement: calculatePlacement(effect, trigger.name), + replication: trigger.name === "failover", + oss_cluster_api_preferred_ip_type: ipType, + oss_cluster_api_preferred_endpoint_type: endpointType, + }, + cluster: { min_nodes: MIN_NODES, actual_nodes: nodeCount }, + oss_cluster_api: { ip_type: ipType, endpoint_type: endpointType }, + description: `Config (${ipType}/${endpointType})`, + }; + return requirement; + }), + })); +} diff --git a/src/app.ts b/src/app.ts index 2de1c94..565b3a7 100644 --- a/src/app.ts +++ b/src/app.ts @@ -11,12 +11,11 @@ import { type SendResult, } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy.ts"; import { executeAction } from "./actions/index.ts"; +import { generateTriggersForEffect } from "./actions/triggers.ts"; import applyDefaultInterceptors from "./default_interceptors/index.ts"; import ProxyStore, { makeId } from "./proxy-store.ts"; import { type ActionRecord, - type ActionTrigger, - type ListActionTriggersResponse, actionIdParamSchema, actionRequestSchema, connectionIdsQuerySchema, @@ -24,11 +23,11 @@ import { encodingSchema, getConfig, interceptorSchema, + type ListActionTriggersResponse, paramSchema, parseBuffer, proxyConfigSchema, slotMigrateEffectSchema, - type SlotMigrateEffect, } from "./util.ts"; const startNewProxy = (config: ProxyConfig) => { @@ -204,11 +203,12 @@ export function createApp(testConfig?: ExtendedProxyConfig) { actionStore.set(actionId, actionRecord); // Execute the action asynchronously + actionRecord.status = "running"; executeAction(type, parameters, proxyStore, config) .then((result) => { actionRecord.status = result.status; - actionRecord.output = "Done"; - actionRecord.error = result.error; + actionRecord.output = result.output ?? "Done"; + actionRecord.error = result.error ?? null; }) .catch((error) => { actionRecord.status = "failed"; @@ -234,73 +234,16 @@ export function createApp(testConfig?: ExtendedProxyConfig) { }); }); - // Hardcoded triggers for each effect - const triggersMap: Record = { - add: [ - { - name: "add-node-trigger-1", - description: "Trigger when a new node is added to the cluster", - requirements: [ - { dbconfig: {}, cluster: { minNodes: 3 }, description: "Requires at least 3 shards and 3 nodes" }, - ], - }, - { - name: "add-node-trigger-2", - description: "Trigger for rebalancing after node addition", - requirements: [ - { dbconfig: {}, cluster: { healthy: true }, description: "Requires replication enabled and healthy cluster" }, - ], - }, - ], - remove: [ - { - name: "remove-node-trigger-1", - description: "Trigger when a node is removed from the cluster", - requirements: [ - { dbconfig: {}, cluster: { minNodes: 2 }, description: "Requires at least 2 shards and 2 nodes" }, - ], - }, - { - name: "remove-node-trigger-2", - description: "Trigger for slot migration before node removal", - requirements: [ - { dbconfig: {}, cluster: { noFailover: true }, description: "Requires persistence and no ongoing failover" }, - ], - }, - ], - "remove-add": [ - { - name: "remove-add-trigger-1", - description: "Trigger for combined remove and add operation", - requirements: [ - { dbconfig: {}, cluster: { minNodes: 3 }, description: "Requires at least 3 shards and 3 nodes" }, - ], - }, - { - name: "remove-add-trigger-2", - description: "Trigger for atomic node replacement", - requirements: [ - { dbconfig: {}, cluster: { quorum: true }, description: "Requires replication and quorum" }, - ], - }, - ], - "slot-shuffle": [ - { - name: "slot-shuffle-trigger-1", - description: "Trigger for redistributing slots across nodes", - requirements: [ - { dbconfig: {}, cluster: { balanced: false }, description: "Requires at least 2 shards and unbalanced cluster" }, - ], - }, - { - name: "slot-shuffle-trigger-2", - description: "Trigger for optimizing slot distribution", - requirements: [ - { dbconfig: {}, cluster: { healthy: true }, description: "Requires auto-balance enabled and healthy cluster" }, - ], - }, - ], - }; + // GET /action - List all submitted actions + app.get("/action", (c) => { + const actions = Array.from(actionStore.values()).map((action) => ({ + job_id: action.id, + action_type: action.type, + status: action.status, + submitted_at: action.submittedAt.toISOString(), + })); + return c.json({ actions }); + }); // GET /slot-migrate - List action triggers for an effect app.get("/slot-migrate", zValidator("query", slotMigrateEffectSchema), (c) => { @@ -309,7 +252,7 @@ export function createApp(testConfig?: ExtendedProxyConfig) { const response: ListActionTriggersResponse = { effect, cluster: { index: 0, nodes: proxyStore.nodeIds.length }, - triggers: triggersMap[effect], + triggers: generateTriggersForEffect(effect, proxyStore.nodeIds.length), }; return c.json(response); diff --git a/src/scenarios.test.ts b/src/scenarios.test.ts deleted file mode 100644 index 1b57cbb..0000000 --- a/src/scenarios.test.ts +++ /dev/null @@ -1,153 +0,0 @@ -import { afterAll, beforeAll, describe, expect, test } from "bun:test"; -import type { SimpleStringReply } from "@redis/client/dist/lib/RESP/types"; -import { createClient } from "redis"; - -import { getFreePortNumber } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy.ts"; -import { createApp } from "./app"; -import createMockRedisServer from "./mock-server"; - -describe("POST /scenarios", () => { - let app: any; - let proxy: any; - let mockRedisServer: any; - let targetPort: number; - - beforeAll(async () => { - const freePort = await getFreePortNumber(); - targetPort = await getFreePortNumber(); - - mockRedisServer = createMockRedisServer(targetPort); - - const testConfig = { - listenPort: [freePort], - listenHost: "127.0.0.1", - targetHost: "127.0.0.1", - targetPort: targetPort, - timeout: 30000, - enableLogging: true, - apiPort: 3001, - }; - - const appInstance = createApp(testConfig); - app = appInstance.app; - proxy = appInstance.proxy; - - await new Promise((resolve) => setTimeout(resolve, 200)); - }); - - afterAll(async () => { - if (proxy) { - await proxy.stop(); - } - if (mockRedisServer) { - mockRedisServer?.stop(true); - } - }); - - test("POST /scenarios with invalid data", async () => { - const res = await app.request("/scenarios", { - method: "POST", - headers: { - "Content-Type": "application/json", - }, - body: JSON.stringify({}), - }); - - expect(res.status).toBe(400); - }); - - test("POST /scenarios with empty responses", async () => { - const res = await app.request("/scenarios", { - method: "POST", - headers: { - "Content-Type": "application/json", - }, - body: JSON.stringify({ responses: [] }), - }); - - expect(res.status).toBe(400); - }); - - test("POST /scenarios with raw encoding", async () => { - const res = await app.request("/scenarios", { - method: "POST", - headers: { - "Content-Type": "application/json", - }, - body: JSON.stringify({ - responses: ["+FIRST\r\n", "+SECOND\r\n", "+THIRD\r\n"], - encoding: "raw", - }), - }); - - expect(res.status).toBe(200); - const result = await res.json(); - expect(result.success).toBe(true); - expect(result.totalResponses).toBe(3); - }); - - test("POST /scenarios with base64 encoding", async () => { - const response1 = Buffer.from("+RESPONSE1\r\n").toString("base64"); - const response2 = Buffer.from("+RESPONSE2\r\n").toString("base64"); - - const res = await app.request("/scenarios", { - method: "POST", - headers: { - "Content-Type": "application/json", - }, - body: JSON.stringify({ - responses: [response1, response2], - encoding: "base64", - }), - }); - - expect(res.status).toBe(200); - const result = await res.json(); - expect(result.success).toBe(true); - expect(result.totalResponses).toBe(2); - }); - - test("Scenario interceptor returns responses sequentially then passes through", async () => { - const client = createClient({ - socket: { - host: "127.0.0.1", - port: proxy.config.listenPort, - }, - }); - - await client.connect(); - await new Promise((resolve) => setTimeout(resolve, 100)); - - // Set up scenario with 2 responses - const scenarioRes = await app.request("/scenarios", { - method: "POST", - headers: { - "Content-Type": "application/json", - }, - body: JSON.stringify({ - responses: ["+SCENARIO1\r\n", "+SCENARIO2\r\n"], - encoding: "raw", - }), - }); - - expect(scenarioRes.status).toBe(200); - - // First command should get first scenario response - const result1 = await client.sendCommand(["PING"]); - expect(result1).toBe("SCENARIO1" as unknown as SimpleStringReply); - - // Second command should get second scenario response - const result2 = await client.sendCommand(["PING"]); - expect(result2).toBe("SCENARIO2" as unknown as SimpleStringReply); - - // Third command should pass through to real server - const result3 = await client.sendCommand(["PING"]); - expect(result3).toBe("PONG" as unknown as SimpleStringReply); - - // Fourth command should also pass through - const result4 = await client.sendCommand(["FOO"]); - expect(result4).toBe("BAR" as unknown as SimpleStringReply); - - await client.disconnect(); - }); -}); diff --git a/src/scenarios/add.ts b/src/scenarios/add.ts deleted file mode 100644 index 960bb2c..0000000 --- a/src/scenarios/add.ts +++ /dev/null @@ -1,119 +0,0 @@ -import type { Context } from "hono"; -import type { ProxyConfig } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy"; -import type ProxyStore from "../proxy-store"; - -import type { ExtendedProxyConfig } from "../util"; -import { - addNode, - buildSMigratedNotification, - buildSMigratingNotification, - createCustomClusterSlotsInterceptor, - findNextAvailablePort, - getSlotRangesForProxy, -} from "./helpers"; -import { getNextSequenceId } from "./sequence-gen"; - -/** - * "ADD NODE" Scenario: - * - * 1. Add a new node to the cluster - * 2. Calculate new slot distribution (slots redistributed from existing nodes to new node) - * 3. Intercept CLUSTER SLOTS to return all nodes (existing + new) with updated slot ranges - * 4. Send SMIGRATING notifications from nodes that will lose slots - * 5. Send SMIGRATED notifications indicating slots moved to the new node - * 6. All nodes remain active - */ -export default async function addNodeScenario( - c: Context, - proxyStore: ProxyStore, - config: ExtendedProxyConfig, -) { - const allProxies = proxyStore.proxies; - if (allProxies.length === 0) { - return c.json( - { - success: false, - error: "No proxies available", - }, - 400, - ); - } - - // Step 1: Get current slot distribution before adding new node - const oldSlotDistribution = allProxies.map((proxy) => ({ - proxy, - slotRanges: getSlotRangesForProxy(proxy, allProxies), - })); - - // Step 2: Add a new node - const newPort = findNextAvailablePort(allProxies); - const newProxyConfig: ProxyConfig = { - ...config, - listenPort: newPort, - }; - - const { nodeId: newNodeId, proxy: newProxy } = addNode(proxyStore, newProxyConfig); - - // Step 3: Calculate new slot distribution with the new node included - const allProxiesWithNew = [...allProxies, newProxy]; - const newSlotDistribution = allProxiesWithNew.map((proxy) => ({ - proxy, - slotRanges: getSlotRangesForProxy(proxy, allProxiesWithNew), - })); - - // Step 4: Intercept CLUSTER SLOTS to return all nodes with new distribution - const clusterSlotsInterceptor = createCustomClusterSlotsInterceptor(allProxiesWithNew); - - for (const proxy of proxyStore.proxies) { - proxy.addGlobalInterceptor(clusterSlotsInterceptor); - } - - // Step 5: Send SMIGRATING notifications from nodes that will lose slots - // Each existing node loses some slots to make room for the new node - for (const { proxy, slotRanges } of oldSlotDistribution) { - const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); - proxy.sendToAllClients(sMigratingBuffer); - } - - // Step 6: Send SMIGRATED notifications after a delay - // This indicates where slots have moved (to the new node) - setTimeout(() => { - const newNodeSlotRanges = getSlotRangesForProxy(newProxy, allProxiesWithNew); - - // Send SMIGRATED from each old node indicating their slots moved to new node - for (const { proxy } of oldSlotDistribution) { - const sMigratedBuffer = buildSMigratedNotification( - [ - { - targetNode: { - host: newProxy.config.listenHost, - port: newProxy.config.listenPort, - }, - slotRanges: newNodeSlotRanges, - }, - ], - getNextSequenceId(), - ); - proxy.sendToAllClients(sMigratedBuffer); - } - }, 5000); - - return c.json({ - success: true, - scenario: "add", - message: "Add node scenario started successfully", - details: { - newNode: `${newProxy.config.listenHost}:${newProxy.config.listenPort}`, - newNodeId, - newNodeSlots: getSlotRangesForProxy(newProxy, allProxiesWithNew), - oldDistribution: oldSlotDistribution.map(({ proxy, slotRanges }) => ({ - node: `${proxy.config.listenHost}:${proxy.config.listenPort}`, - slots: slotRanges, - })), - newDistribution: newSlotDistribution.map(({ proxy, slotRanges }) => ({ - node: `${proxy.config.listenHost}:${proxy.config.listenPort}`, - slots: slotRanges, - })), - }, - }); -} diff --git a/src/scenarios/index.ts b/src/scenarios/index.ts deleted file mode 100644 index db5e4c2..0000000 --- a/src/scenarios/index.ts +++ /dev/null @@ -1,37 +0,0 @@ -import type { Context } from "hono"; -import type { z } from "zod"; -import type ProxyStore from "../proxy-store"; -import type { - ExtendedProxyConfig, - predefinedScenarioParamSchema, -} from "../util"; -import addNodeScenario from "./add"; -import removeSrcAddDestScenario from "./remove-add"; -import removeNodeScenario from "./remove"; -import slotShuffleScenario from "./slot-shuffle"; - -type PredefinedScenario = z.infer< - typeof predefinedScenarioParamSchema ->["scenario"]; - -export default async function applyPredefinedScenario( - scenario: PredefinedScenario, - c: Context, - proxyStore: ProxyStore, - config: ExtendedProxyConfig, -) { - switch (scenario) { - case "remove-add": - return await removeSrcAddDestScenario(c, proxyStore, config); - case "add": - return await addNodeScenario(c, proxyStore, config); - case "remove": - return await removeNodeScenario(c, proxyStore, config); - case "slot-shuffle": - return await slotShuffleScenario(c, proxyStore, config); - - default: - // This should never happen due to Zod validation, but TypeScript requires it - return c.json({ success: false, error: "Unknown scenario" }, 400); - } -} diff --git a/src/scenarios/remove-add.ts b/src/scenarios/remove-add.ts deleted file mode 100644 index a2b7570..0000000 --- a/src/scenarios/remove-add.ts +++ /dev/null @@ -1,121 +0,0 @@ -import type { Context } from "hono"; -import type { ProxyConfig } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy"; -import type ProxyStore from "../proxy-store"; - -import { makeId } from "../proxy-store"; -import type { ExtendedProxyConfig } from "../util"; -import { - addNode, - buildSMigratedNotification, - buildSMigratingNotification, - createCustomClusterSlotsInterceptor, - findNextAvailablePort, - getSlotRangesForProxy, - pickRandom, -} from "./helpers"; -import { getNextSequenceId } from "./sequence-gen"; - -/** - * "REMOVED AND ADDED" Scenario: - * - * 1. Pick a random node from existing proxies - * 2. Add one more node - * 3. Intercept CLUSTER SLOTS to return all nodes + the new one - the randomly selected one - * 4. Send SMIGRATING notification to all clients (slots about to be migrated) - * 5. Send SMIGRATED notification to all clients (slots migrated from picked node to new node) - * 6. Kill the old node - */ -export default async function removeSrcAddDestScenario( - c: Context, - proxyStore: ProxyStore, - config: ExtendedProxyConfig, -) { - // Step 1: Pick a random node - const allProxies = proxyStore.proxies; - if (allProxies.length === 0) { - return c.json( - { - success: false, - error: "No proxies available to select from", - }, - 400, - ); - } - - const proxyToBeRemoved = pickRandom(allProxies); - if (!proxyToBeRemoved) { - return c.json( - { - success: false, - error: "Failed to select a random proxy", - }, - 400, - ); - } - - // Get the slot ranges for the randomly selected proxy before adding the new node - const slotRanges = getSlotRangesForProxy(proxyToBeRemoved, allProxies); - - // Step 2: Add one more node - const newPort = findNextAvailablePort(allProxies); - const newProxyConfig: ProxyConfig = { - ...config, - listenPort: newPort, - }; - - const { nodeId: newNodeId, proxy: newProxy } = addNode(proxyStore, newProxyConfig); - - // Step 3: Create list of proxies excluding the randomly selected one, and including the new one - const proxiesForClusterSlots = allProxies.filter((p) => p !== proxyToBeRemoved).concat(newProxy); - - // Add the custom cluster slots interceptor to all proxies - const clusterSlotsInterceptor = createCustomClusterSlotsInterceptor(proxiesForClusterSlots); - - for (const proxy of proxyStore.proxies) { - proxy.addGlobalInterceptor(clusterSlotsInterceptor); - } - - // Step 4: Send SMIGRATING notification to all clients - // This notifies clients that slots are about to be migrated - const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); - proxyToBeRemoved.sendToAllClients(sMigratingBuffer); - - // Step 5: Send SMIGRATED notification to all clients - // This notifies clients that slots from the picked node have moved to the new node - setTimeout(() => { - const sMigratedBuffer = buildSMigratedNotification( - [ - { - targetNode: { - host: newProxy.config.listenHost, - port: newProxy.config.listenPort, - }, - slotRanges, - }, - ], - getNextSequenceId(), - ); - proxyToBeRemoved.sendToAllClients(sMigratedBuffer); - - setTimeout(() => { - const { targetHost, targetPort, listenPort } = proxyToBeRemoved.config; - proxyStore.delete(makeId(targetHost, targetPort, listenPort)); - }, 2000); - }, 5000); - - return c.json({ - success: true, - scenario: "foo", - message: "Foo scenario started successfully", - details: { - excludedNode: `${proxyToBeRemoved.config.listenHost}:${proxyToBeRemoved.config.listenPort}`, - newNode: `${newProxy.config.listenHost}:${newProxy.config.listenPort}`, - newNodeId, - migratedSlots: slotRanges, - clusterSlotsProxies: proxiesForClusterSlots.map((p) => ({ - host: p.config.listenHost, - port: p.config.listenPort, - })), - }, - }); -} diff --git a/src/scenarios/remove.ts b/src/scenarios/remove.ts deleted file mode 100644 index da02c1e..0000000 --- a/src/scenarios/remove.ts +++ /dev/null @@ -1,125 +0,0 @@ -import type { Context } from "hono"; -import type ProxyStore from "../proxy-store"; - -import { makeId } from "../proxy-store"; -import type { ExtendedProxyConfig } from "../util"; -import { - buildSMigratedNotification, - buildSMigratingNotification, - createCustomClusterSlotsInterceptor, - getSlotRangesForProxy, - pickRandom, -} from "./helpers"; -import { getNextSequenceId } from "./sequence-gen"; - -/** - * "REMOVE NODE" Scenario: - * - * 1. Pick a random node from existing proxies to remove - * 2. Get the slot ranges currently assigned to that node - * 3. Calculate new slot distribution among remaining nodes - * 4. Intercept CLUSTER SLOTS to return all nodes except the one being removed - * 5. Send SMIGRATING notification from the node being removed - * 6. Send SMIGRATED notifications indicating where slots moved (to remaining nodes) - * 7. Kill/delete the removed node after notifications - */ -export default async function removeNodeScenario( - c: Context, - proxyStore: ProxyStore, - _config: ExtendedProxyConfig, -) { - // Step 1: Pick a random node to remove - const allProxies = proxyStore.proxies; - if (allProxies.length === 0) { - return c.json( - { - success: false, - error: "No proxies available to select from", - }, - 400, - ); - } - - if (allProxies.length === 1) { - return c.json( - { - success: false, - error: "Cannot remove the last remaining node", - }, - 400, - ); - } - - const proxyToBeRemoved = pickRandom(allProxies); - if (!proxyToBeRemoved) { - return c.json( - { - success: false, - error: "Failed to select a random proxy", - }, - 400, - ); - } - - // Step 2: Get the slot ranges for the node being removed - const removedNodeSlotRanges = getSlotRangesForProxy(proxyToBeRemoved, allProxies); - - // Step 3: Create list of remaining proxies (excluding the one being removed) - const remainingProxies = allProxies.filter((p) => p !== proxyToBeRemoved); - - // Calculate new slot distribution for remaining nodes - const newSlotDistribution = remainingProxies.map((proxy) => ({ - proxy, - slotRanges: getSlotRangesForProxy(proxy, remainingProxies), - })); - - // Step 4: Intercept CLUSTER SLOTS to return only remaining nodes - const clusterSlotsInterceptor = createCustomClusterSlotsInterceptor(remainingProxies); - - for (const proxy of proxyStore.proxies) { - proxy.addGlobalInterceptor(clusterSlotsInterceptor); - } - - // Step 5: Send SMIGRATING notification from the node being removed - const sMigratingBuffer = buildSMigratingNotification(removedNodeSlotRanges, getNextSequenceId()); - proxyToBeRemoved.sendToAllClients(sMigratingBuffer); - - // Step 6: Send SMIGRATED notifications after a delay - // This indicates where slots from the removed node have moved (distributed to remaining nodes) - setTimeout(() => { - // Build SMIGRATED notification showing slots redistributed to remaining nodes - const migratedSlots = newSlotDistribution.map(({ proxy, slotRanges }) => ({ - targetNode: { - host: proxy.config.listenHost, - port: proxy.config.listenPort, - }, - slotRanges, - })); - - const sMigratedBuffer = buildSMigratedNotification(migratedSlots, getNextSequenceId()); - proxyToBeRemoved.sendToAllClients(sMigratedBuffer); - - // Step 7: Kill the removed node after another delay - setTimeout(() => { - const { targetHost, targetPort, listenPort } = proxyToBeRemoved.config; - proxyStore.delete(makeId(targetHost, targetPort, listenPort)); - }, 2000); - }, 5000); - - return c.json({ - success: true, - scenario: "remove", - message: "Remove node scenario started successfully", - details: { - removedNode: `${proxyToBeRemoved.config.listenHost}:${proxyToBeRemoved.config.listenPort}`, - removedNodeSlots: removedNodeSlotRanges, - remainingNodes: remainingProxies.map((p) => ({ - node: `${p.config.listenHost}:${p.config.listenPort}`, - })), - newDistribution: newSlotDistribution.map(({ proxy, slotRanges }) => ({ - node: `${proxy.config.listenHost}:${proxy.config.listenPort}`, - slots: slotRanges, - })), - }, - }); -} diff --git a/src/scenarios/slot-shuffle.ts b/src/scenarios/slot-shuffle.ts deleted file mode 100644 index ab06138..0000000 --- a/src/scenarios/slot-shuffle.ts +++ /dev/null @@ -1,150 +0,0 @@ -import type { Context } from "hono"; -import type { RedisProxy } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy"; -import type ProxyStore from "../proxy-store"; - -import type { ExtendedProxyConfig } from "../util"; -import { - buildSMigratedNotification, - buildSMigratingNotification, - getSlotRangesForProxy, -} from "./helpers"; -import { getNextSequenceId } from "./sequence-gen"; - -/** - * "SLOT SHUFFLE" Scenario: - * - * 1. Keep all existing nodes (no add/remove) - * 2. Randomly shuffle slot assignments among existing nodes - * 3. Intercept CLUSTER SLOTS to return the new shuffled slot distribution - * 4. Send SMIGRATING notifications from nodes losing slots - * 5. Send SMIGRATED notifications indicating the new slot ownership - * 6. All nodes remain active with new slot assignments - */ -export default async function slotShuffleScenario( - c: Context, - proxyStore: ProxyStore, - _config: ExtendedProxyConfig, -) { - const allProxies = proxyStore.proxies; - if (allProxies.length === 0) { - return c.json( - { - success: false, - error: "No proxies available", - }, - 400, - ); - } - - if (allProxies.length === 1) { - return c.json( - { - success: false, - error: "Cannot shuffle slots with only one node", - }, - 400, - ); - } - - // Step 1: Get current slot distribution - const oldSlotDistribution = allProxies.map((proxy) => ({ - proxy, - slotRanges: getSlotRangesForProxy(proxy, allProxies), - })); - - // Step 2: Shuffle the proxies array to create a new slot distribution - // We'll create a shuffled copy of the proxies array - const shuffledProxies = [...allProxies]; - - // Fisher-Yates shuffle algorithm - for (let i = shuffledProxies.length - 1; i > 0; i--) { - const j = Math.floor(Math.random() * (i + 1)); - const temp = shuffledProxies[i]; - const jProxy = shuffledProxies[j]; - if (temp && jProxy) { - shuffledProxies[i] = jProxy; - shuffledProxies[j] = temp; - } - } - - // Step 3: Calculate new slot distribution based on shuffled order - // We'll use the shuffled order to reassign slots - const newSlotDistribution = shuffledProxies.map((proxy, index) => { - const slotLength = Math.floor(16384 / shuffledProxies.length); - const from = index * slotLength; - const to = index === shuffledProxies.length - 1 ? 16383 : from + slotLength - 1; - return { - proxy, - slotRanges: `${from}-${to}`, - }; - }); - - // Step 4: Create a custom interceptor that returns the shuffled slot distribution - // We need to create a custom interceptor because the slots don't match the original order - const customShuffledInterceptor = { - name: "cluster-simulation-interceptor", - fn: async (data: Buffer, next: any, state: any) => { - state.invokeCount++; - - if (data.toString().toLowerCase() !== "*2\r\n$7\r\ncluster\r\n$5\r\nslots\r\n") { - return next(data); - } - - state.matchCount++; - - const mapping = newSlotDistribution.map(({ proxy, slotRanges }) => { - const [from, to] = slotRanges.split("-").map(Number); - const id = `proxy-id-${proxy.config.listenPort}`; - return `*3\r\n:${from}\r\n:${to}\r\n*3\r\n$${proxy.config.listenHost.length}\r\n${proxy.config.listenHost}\r\n:${proxy.config.listenPort}\r\n$${id.length}\r\n${id}\r\n`; - }); - - const response = `*${newSlotDistribution.length}\r\n${mapping.join("")}`; - return Buffer.from(response); - }, - }; - - for (const proxy of proxyStore.proxies) { - proxy.addGlobalInterceptor(customShuffledInterceptor); - } - - // Step 5: Send SMIGRATING notifications from all nodes - // All nodes are potentially losing and gaining slots - for (const { proxy, slotRanges } of oldSlotDistribution) { - const sMigratingBuffer = buildSMigratingNotification(slotRanges, getNextSequenceId()); - proxy.sendToAllClients(sMigratingBuffer); - } - - // Step 6: Send SMIGRATED notifications after a delay - // This indicates the new slot ownership across all nodes - setTimeout(() => { - // Each node sends SMIGRATED showing where all slots now live - for (const { proxy: sourceProxy } of oldSlotDistribution) { - const migratedSlots = newSlotDistribution.map(({ proxy, slotRanges }) => ({ - targetNode: { - host: proxy.config.listenHost, - port: proxy.config.listenPort, - }, - slotRanges, - })); - - const sMigratedBuffer = buildSMigratedNotification(migratedSlots, getNextSequenceId()); - sourceProxy.sendToAllClients(sMigratedBuffer); - } - }, 5000); - - return c.json({ - success: true, - scenario: "slot-shuffle", - message: "Slot shuffle scenario started successfully", - details: { - oldDistribution: oldSlotDistribution.map(({ proxy, slotRanges }) => ({ - node: `${proxy.config.listenHost}:${proxy.config.listenPort}`, - slots: slotRanges, - })), - newDistribution: newSlotDistribution.map(({ proxy, slotRanges }) => ({ - node: `${proxy.config.listenHost}:${proxy.config.listenPort}`, - slots: slotRanges, - })), - }, - }); -} diff --git a/src/util.ts b/src/util.ts index 203f45e..0ca5d74 100644 --- a/src/util.ts +++ b/src/util.ts @@ -31,10 +31,6 @@ export const interceptorSchema = z.object({ response: z.string(), }); -export const predefinedScenarioParamSchema = z.object({ - scenario: z.enum(["remove-add", "remove", "add", "slot-shuffle"]), -}); - export const slotMigrateEffectSchema = z.object({ effect: z.enum(["add", "remove", "remove-add", "slot-shuffle"]), }); @@ -42,8 +38,9 @@ export const slotMigrateEffectSchema = z.object({ export type SlotMigrateEffect = z.infer["effect"]; export interface ActionTriggerRequirement { - dbconfig: unknown; - cluster: unknown; + dbconfig: Record & { name: string }; + cluster: { min_nodes: number; actual_nodes: number }; + oss_cluster_api: { ip_type: string; endpoint_type: string }; description: string; } @@ -66,11 +63,13 @@ export const actionTypeSchema = z.enum([ "reshard", "sequence_of_actions", "network_failure", + "network_latency", "execute_rlutil_command", "execute_rladmin_command", "enable_entraid", "upgrade", "wait", + "wait_for_database_active", "migrate", "bind", "update_cluster_config", @@ -82,6 +81,9 @@ export const actionTypeSchema = z.enum([ "proxy_failure", "cluster_failure", "slot_migrate", + "topology_change_standalone", + "reset_cluster", + "collect_debuginfo", ]); export type ActionType = z.infer; From 572aebeddcbf62c2bf655355af90350efcbe4102 Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Thu, 10 Sep 2026 16:56:09 +0300 Subject: [PATCH 5/5] feat: add reject-traffic endpoints to simulate an offline endpoint POST /reject-traffic/start drops every client connection and stops accepting new ones on all nodes. POST /reject-traffic/stop brings the listeners back up. Interceptors and topology survive the cycle, and both endpoints are idempotent. Co-Authored-By: Claude Fable 5 --- src/app.ts | 24 ++++++++ src/reject-traffic.test.ts | 117 +++++++++++++++++++++++++++++++++++++ 2 files changed, 141 insertions(+) create mode 100644 src/reject-traffic.test.ts diff --git a/src/app.ts b/src/app.ts index 565b3a7..00851e1 100644 --- a/src/app.ts +++ b/src/app.ts @@ -51,6 +51,30 @@ export function createApp(testConfig?: ExtendedProxyConfig) { config.defaultInterceptors && applyDefaultInterceptors(config.defaultInterceptors, proxyStore); + // Simulate the endpoints being offline: drop every client connection and + // stop accepting new ones until rejecting is stopped again. + let rejectingTraffic = false; + + app.post("/reject-traffic/start", async (c) => { + if (!rejectingTraffic) { + rejectingTraffic = true; + for (const proxy of proxyStore.proxies) { + await proxy.stop(); + } + } + return c.json({ success: true, rejecting: rejectingTraffic }); + }); + + app.post("/reject-traffic/stop", async (c) => { + if (rejectingTraffic) { + rejectingTraffic = false; + for (const proxy of proxyStore.proxies) { + await proxy.start(); + } + } + return c.json({ success: true, rejecting: rejectingTraffic }); + }); + app.post("/nodes", zValidator("json", proxyConfigSchema), async (c) => { const data = await c.req.json(); const cfg: ProxyConfig = { ...config, ...data }; diff --git a/src/reject-traffic.test.ts b/src/reject-traffic.test.ts new file mode 100644 index 0000000..a1bac6e --- /dev/null +++ b/src/reject-traffic.test.ts @@ -0,0 +1,117 @@ +import { afterAll, beforeAll, describe, expect, test } from "bun:test"; +import type { Socket } from "bun"; + +import { getFreePortNumber } from "redis-monorepo/packages/test-utils/lib/proxy/redis-proxy.ts"; +import { createApp } from "./app"; +import createMockRedisServer from "./mock-server"; + +describe("Reject traffic", () => { + let app: any; + let mockRedisServer: ReturnType; + let listenPort: number; + let targetPort: number; + + const connectClient = async (port: number) => { + let socket: Socket | undefined; + let closed = false; + let onClose = () => {}; + await new Promise((resolve, reject) => { + Bun.connect({ + hostname: "127.0.0.1", + port, + socket: { + open(openedSocket) { + socket = openedSocket; + resolve(); + }, + data() {}, + error(_socket, error) { + reject(error); + }, + close() { + closed = true; + onClose(); + }, + }, + }).catch(reject); + }); + return { + isClosed: () => closed, + waitForClose: () => + new Promise((resolve, reject) => { + if (closed) return resolve(); + onClose = resolve; + setTimeout(() => reject(new Error("connection was not closed")), 2000); + }), + close: () => socket?.end(), + }; + }; + + beforeAll(async () => { + listenPort = await getFreePortNumber(); + targetPort = await getFreePortNumber(); + + mockRedisServer = createMockRedisServer(targetPort); + + const appInstance = createApp({ + listenPort: [listenPort], + listenHost: "127.0.0.1", + targetHost: "127.0.0.1", + targetPort: targetPort, + timeout: 30000, + enableLogging: false, + apiPort: 3003, + }); + app = appInstance.app; + + await new Promise((resolve) => setTimeout(resolve, 200)); + }); + + afterAll(async () => { + const res = await app.request("/nodes"); + const { ids } = await res.json(); + for (const id of ids) { + await app.request(`/nodes/${encodeURIComponent(id)}`, { method: "DELETE" }); + } + mockRedisServer?.stop(true); + }); + + test("start drops connections and refuses new ones, stop restores service", async () => { + const client = await connectClient(listenPort); + expect(client.isClosed()).toBe(false); + + const startRes = await app.request("/reject-traffic/start", { method: "POST" }); + expect(startRes.status).toBe(200); + expect(await startRes.json()).toEqual({ success: true, rejecting: true }); + + // The existing connection is dropped + await client.waitForClose(); + + // New connections are refused while rejecting + await expect(connectClient(listenPort)).rejects.toThrow(); + + const stopRes = await app.request("/reject-traffic/stop", { method: "POST" }); + expect(stopRes.status).toBe(200); + expect(await stopRes.json()).toEqual({ success: true, rejecting: false }); + + // Service is back: new connections are accepted again + const revivedClient = await connectClient(listenPort); + expect(revivedClient.isClosed()).toBe(false); + revivedClient.close(); + }); + + test("start and stop are idempotent", async () => { + for (const _ of [1, 2]) { + const res = await app.request("/reject-traffic/start", { method: "POST" }); + expect((await res.json()).rejecting).toBe(true); + } + for (const _ of [1, 2]) { + const res = await app.request("/reject-traffic/stop", { method: "POST" }); + expect((await res.json()).rejecting).toBe(false); + } + + const client = await connectClient(listenPort); + expect(client.isClosed()).toBe(false); + client.close(); + }); +});