diff --git a/fluxer_api/pkgs/kv_client/package.json b/fluxer_api/pkgs/kv_client/package.json index 888c9586b..3a19c30b7 100644 --- a/fluxer_api/pkgs/kv_client/package.json +++ b/fluxer_api/pkgs/kv_client/package.json @@ -7,6 +7,8 @@ "./*": "./*" }, "scripts": { + "test": "vitest run", + "test:watch": "vitest", "typecheck": "tsgo --noEmit" }, "dependencies": { @@ -16,6 +18,8 @@ }, "devDependencies": { "@types/node": "catalog:", - "@typescript/native-preview": "catalog:" + "@typescript/native-preview": "catalog:", + "vite-tsconfig-paths": "catalog:", + "vitest": "catalog:" } } diff --git a/fluxer_api/pkgs/kv_client/src/KVClient.ts b/fluxer_api/pkgs/kv_client/src/KVClient.ts index ce8012723..2fe48d739 100644 --- a/fluxer_api/pkgs/kv_client/src/KVClient.ts +++ b/fluxer_api/pkgs/kv_client/src/KVClient.ts @@ -1,5 +1,6 @@ // SPDX-License-Identifier: AGPL-3.0-or-later +import {createHash} from 'node:crypto'; import type {IKVPipeline, IKVProvider, IKVSubscription, KVRateLimitResult} from '@pkgs/kv_client/src/IKVProvider'; import { type IKVLogger, @@ -232,6 +233,8 @@ redis.call('DEL', KEYS[2]) return 1 `; +const SCRIPT_SHA_CACHE = new Map(); + interface ScriptPurgeBatchResult { urls: Array; tokens: number; @@ -699,7 +702,18 @@ export class KVClient implements IKVProvider { keyCount: number, ...args: Array ): Promise { - return await this.execute(command, async () => this.client.eval(script, keyCount, ...args)); + return await this.execute(command, async () => this.evalCachedScript(script, keyCount, args)); + } + + private async evalCachedScript(script: string, keyCount: number, args: Array): Promise { + try { + return await this.client.evalsha(getScriptSha(script), keyCount, ...args); + } catch (error) { + if (!isNoScriptError(error)) { + throw error; + } + return await this.client.eval(script, keyCount, ...args); + } } private async executeJsonScript( @@ -795,6 +809,23 @@ function normalizeRateLimitResult(result: KVRateLimitResult): KVRateLimitResult }; } +function getScriptSha(script: string): string { + const cached = SCRIPT_SHA_CACHE.get(script); + if (cached !== undefined) { + return cached; + } + const sha = createHash('sha1').update(script).digest('hex'); + SCRIPT_SHA_CACHE.set(script, sha); + return sha; +} + +function isNoScriptError(error: unknown): boolean { + if (!(error instanceof Error)) { + return false; + } + return error.message.includes('NOSCRIPT'); +} + function isTimeoutError(error: unknown): boolean { if (!(error instanceof Error)) { return false; diff --git a/fluxer_api/pkgs/kv_client/src/__tests__/KVClient.test.ts b/fluxer_api/pkgs/kv_client/src/__tests__/KVClient.test.ts new file mode 100644 index 000000000..ec1f8c97b --- /dev/null +++ b/fluxer_api/pkgs/kv_client/src/__tests__/KVClient.test.ts @@ -0,0 +1,209 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {createHash} from 'node:crypto'; +import {KVClient} from '@pkgs/kv_client/src/KVClient'; +import {KVClientErrorCode} from '@pkgs/kv_client/src/KVClientError'; +import {beforeEach, describe, expect, it, vi} from 'vitest'; + +const {evalshaMock, evalMock} = vi.hoisted(() => ({ + evalshaMock: vi.fn(), + evalMock: vi.fn(), +})); + +vi.mock('ioredis', () => { + class MockRedis { + evalsha = evalshaMock; + eval = evalMock; + } + return {default: MockRedis, Cluster: MockRedis}; +}); + +const RATE_LIMIT_REPLY = JSON.stringify({ + allowed: true, + limit: 5, + remaining: 4, + resetAfterMs: 200, + resetAtMs: 1717171717, + retryAfterMs: 0, +}); + +const EXPECTED_RATE_LIMIT_RESULT = { + allowed: true, + limit: 5, + remaining: 4, + resetAfterMs: 200, + resetAtMs: 1717171717, + retryAfterMs: 0, +}; + +function createClient(): KVClient { + return new KVClient('redis://127.0.0.1:6379'); +} + +function noScriptError(): Error { + return new Error('NOSCRIPT No matching script. Please use EVAL.'); +} + +function getCallArguments(mock: typeof evalshaMock, index: number): Array { + const call = mock.mock.calls[index]; + if (!call) { + throw new Error(`Expected a call at index ${index}`); + } + return call; +} + +describe('KVClient script execution', () => { + beforeEach(() => { + evalshaMock.mockReset(); + evalMock.mockReset(); + }); + + it('sends EVALSHA instead of EVAL for the leaky bucket rate limit script', async () => { + evalshaMock.mockResolvedValue(RATE_LIMIT_REPLY); + const result = await createClient().checkLeakyBucketLimit('rate_limit:bucket', 5, 1000, 1); + expect(evalMock).not.toHaveBeenCalled(); + expect(evalshaMock).toHaveBeenCalledTimes(1); + const [sha, keyCount, key, , limit, windowMs, cost] = getCallArguments(evalshaMock, 0); + expect(sha).toMatch(/^[0-9a-f]{40}$/); + expect(keyCount).toBe(1); + expect(key).toBe('rate_limit:bucket'); + expect(limit).toBe(5); + expect(windowMs).toBe(1000); + expect(cost).toBe(1); + expect(result).toEqual(EXPECTED_RATE_LIMIT_RESULT); + }); + + it('falls back to EVAL with the original script and identical arguments on NOSCRIPT', async () => { + evalshaMock.mockRejectedValue(noScriptError()); + evalMock.mockResolvedValue(RATE_LIMIT_REPLY); + const result = await createClient().checkLeakyBucketLimit('rate_limit:bucket', 5, 1000, 1); + expect(evalshaMock).toHaveBeenCalledTimes(1); + expect(evalMock).toHaveBeenCalledTimes(1); + const [sha, ...evalshaRest] = getCallArguments(evalshaMock, 0); + const [script, ...evalRest] = getCallArguments(evalMock, 0); + expect(createHash('sha1').update(String(script)).digest('hex')).toBe(sha); + expect(evalRest).toEqual(evalshaRest); + expect(String(script)).toContain("local rawState = redis.call('GET', key)"); + expect(result).toEqual(EXPECTED_RATE_LIMIT_RESULT); + }); + + it('keeps using EVALSHA with the same digest after a NOSCRIPT fallback', async () => { + evalshaMock.mockRejectedValueOnce(noScriptError()).mockResolvedValue(RATE_LIMIT_REPLY); + evalMock.mockResolvedValue(RATE_LIMIT_REPLY); + const client = createClient(); + const first = await client.checkLeakyBucketLimit('rate_limit:bucket', 5, 1000, 1); + const second = await client.checkLeakyBucketLimit('rate_limit:bucket', 5, 1000, 1); + expect(evalshaMock).toHaveBeenCalledTimes(2); + expect(evalMock).toHaveBeenCalledTimes(1); + expect(getCallArguments(evalshaMock, 1)[0]).toBe(getCallArguments(evalshaMock, 0)[0]); + expect(first).toEqual(EXPECTED_RATE_LIMIT_RESULT); + expect(second).toEqual(EXPECTED_RATE_LIMIT_RESULT); + }); + + it('sends EVALSHA for every scripted command', async () => { + const cases: Array<{name: string; reply: unknown; keyCount: number; run: (client: KVClient) => Promise}> = + [ + { + name: 'releaseLock', + reply: 1, + keyCount: 1, + run: async (client) => client.releaseLock('lock:key', 'token'), + }, + { + name: 'extendLock', + reply: 1, + keyCount: 1, + run: async (client) => client.extendLock('lock:key', 'token', 30), + }, + { + name: 'renewSnowflakeNode', + reply: 1, + keyCount: 1, + run: async (client) => client.renewSnowflakeNode('snowflake:1', 'instance', 30), + }, + { + name: 'checkLeakyBucketLimit', + reply: RATE_LIMIT_REPLY, + keyCount: 1, + run: async (client) => client.checkLeakyBucketLimit('rate_limit:bucket', 5, 1000, 1), + }, + { + name: 'tryConsumeTokens', + reply: 2, + keyCount: 1, + run: async (client) => client.tryConsumeTokens('tokens:key', 2, 10, 1, 1000), + }, + { + name: 'scheduleBulkDeletion', + reply: 1, + keyCount: 2, + run: async (client) => client.scheduleBulkDeletion('queue:key', 'secondary:key', 1, 'value'), + }, + { + name: 'removeBulkDeletion', + reply: 1, + keyCount: 2, + run: async (client) => client.removeBulkDeletion('queue:key', 'secondary:key'), + }, + { + name: 'dequeuePurgeBatch', + reply: JSON.stringify({urls: ['https://fluxer.test/a.png'], tokens: 1}), + keyCount: 2, + run: async (client) => client.dequeuePurgeBatch('queue:key', 'bucket:key', 10, 10, 1, 1000), + }, + { + name: 'evalScript', + reply: 1, + keyCount: 1, + run: async (client) => client.evalScript('customScript', "return redis.call('GET', KEYS[1])", 1, 'key'), + }, + ]; + const digests = new Set(); + for (const scriptCase of cases) { + evalshaMock.mockReset(); + evalMock.mockReset(); + evalshaMock.mockResolvedValue(scriptCase.reply); + await scriptCase.run(createClient()); + expect(evalMock, scriptCase.name).not.toHaveBeenCalled(); + expect(evalshaMock, scriptCase.name).toHaveBeenCalledTimes(1); + const [sha, keyCount] = getCallArguments(evalshaMock, 0); + expect(sha, scriptCase.name).toMatch(/^[0-9a-f]{40}$/); + expect(keyCount, scriptCase.name).toBe(scriptCase.keyCount); + digests.add(sha); + } + expect(digests.size).toBe(cases.length); + }); + + it('does not retry with EVAL when the script fails for another reason', async () => { + evalshaMock.mockRejectedValue(new Error('Connection is closed.')); + await expect(createClient().releaseLock('lock:key', 'token')).rejects.toMatchObject({ + code: KVClientErrorCode.REQUEST_FAILED, + message: 'KV request failed (releaseLock): Connection is closed.', + }); + expect(evalMock).not.toHaveBeenCalled(); + }); + + it('normalizes timeouts raised by the EVALSHA attempt', async () => { + evalshaMock.mockRejectedValue(new Error('Command timed out')); + await expect(createClient().releaseLock('lock:key', 'token')).rejects.toMatchObject({ + code: KVClientErrorCode.TIMEOUT, + message: 'KV request timed out: releaseLock', + }); + }); + + it('normalizes failures raised by the EVAL fallback', async () => { + evalshaMock.mockRejectedValue(noScriptError()); + evalMock.mockRejectedValue(new Error('Connection is closed.')); + await expect(createClient().releaseLock('lock:key', 'token')).rejects.toMatchObject({ + code: KVClientErrorCode.REQUEST_FAILED, + message: 'KV request failed (releaseLock): Connection is closed.', + }); + }); + + it('reports invalid JSON from a scripted command', async () => { + evalshaMock.mockResolvedValue('not json'); + await expect(createClient().checkLeakyBucketLimit('rate_limit:bucket', 5, 1000, 1)).rejects.toMatchObject({ + code: KVClientErrorCode.INVALID_RESPONSE, + }); + }); +}); diff --git a/fluxer_api/pkgs/kv_client/vitest.config.ts b/fluxer_api/pkgs/kv_client/vitest.config.ts new file mode 100644 index 000000000..08796abc6 --- /dev/null +++ b/fluxer_api/pkgs/kv_client/vitest.config.ts @@ -0,0 +1,27 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import path from 'node:path'; +import {fileURLToPath} from 'node:url'; +import tsconfigPaths from 'vite-tsconfig-paths'; +import {defineConfig} from 'vitest/config'; + +const __dirname = path.dirname(fileURLToPath(import.meta.url)); + +export default defineConfig({ + plugins: [ + tsconfigPaths({ + root: path.resolve(__dirname, '../..'), + }), + ], + test: { + globals: true, + environment: 'node', + include: ['**/*.{test,spec}.{ts,tsx}'], + exclude: ['node_modules', 'dist'], + coverage: { + provider: 'v8', + reporter: ['text', 'json', 'html'], + exclude: ['**/*.test.tsx', '**/*.spec.tsx', 'node_modules/'], + }, + }, +}); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 97cb49d5f..792f2871a 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -797,6 +797,12 @@ importers: '@typescript/native-preview': specifier: 'catalog:' version: 7.0.0-dev.20260224.1 + vite-tsconfig-paths: + specifier: 'catalog:' + version: 6.1.1(typescript@5.9.3)(vite@7.3.1(@types/node@25.3.0)(jiti@2.6.1)(lightningcss@1.31.1)(tsx@4.21.0)(yaml@2.8.2)) + vitest: + specifier: 'catalog:' + version: 4.0.18(@opentelemetry/api@1.9.0)(@types/node@25.3.0)(@vitest/browser-playwright@4.0.18)(happy-dom@20.7.0)(jiti@2.6.1)(jsdom@28.1.0)(lightningcss@1.31.1)(msw@2.12.10(@types/node@25.3.0)(typescript@5.9.3))(tsx@4.21.0)(yaml@2.8.2) fluxer_api/pkgs/locale: dependencies: