mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
perf(kv-client): use EVALSHA for rate limit scripts (#2130)
This commit is contained in:
@@ -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:"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<string, string>();
|
||||
|
||||
interface ScriptPurgeBatchResult {
|
||||
urls: Array<string>;
|
||||
tokens: number;
|
||||
@@ -699,7 +702,18 @@ export class KVClient implements IKVProvider {
|
||||
keyCount: number,
|
||||
...args: Array<string | number>
|
||||
): Promise<unknown> {
|
||||
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<string | number>): Promise<unknown> {
|
||||
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<T>(
|
||||
@@ -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;
|
||||
|
||||
@@ -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<unknown> {
|
||||
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<unknown>}> =
|
||||
[
|
||||
{
|
||||
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<unknown>();
|
||||
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,
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -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/'],
|
||||
},
|
||||
},
|
||||
});
|
||||
Generated
+6
@@ -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([email protected])([email protected](@types/[email protected])([email protected])([email protected])([email protected])([email protected]))
|
||||
vitest:
|
||||
specifier: 'catalog:'
|
||||
version: 4.0.18(@opentelemetry/[email protected])(@types/[email protected])(@vitest/[email protected])([email protected])([email protected])([email protected])([email protected])([email protected](@types/[email protected])([email protected]))([email protected])([email protected])
|
||||
|
||||
fluxer_api/pkgs/locale:
|
||||
dependencies:
|
||||
|
||||
Reference in New Issue
Block a user