mirror of
https://github.com/fluxerapp/fluxer.git
synced 2026-09-02 21:04:06 +03:00
fix(cache): single-flight getOrSet and cache null results (#2199)
This commit is contained in:
Vendored
+5
-1
@@ -7,6 +7,8 @@
|
||||
"./*": "./*"
|
||||
},
|
||||
"scripts": {
|
||||
"test": "vitest run",
|
||||
"test:watch": "vitest",
|
||||
"typecheck": "tsgo --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
@@ -14,6 +16,8 @@
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "catalog:",
|
||||
"@typescript/native-preview": "catalog:"
|
||||
"@typescript/native-preview": "catalog:",
|
||||
"vite-tsconfig-paths": "catalog:",
|
||||
"vitest": "catalog:"
|
||||
}
|
||||
}
|
||||
|
||||
+9
-3
@@ -1,20 +1,26 @@
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
import type {CacheLogger} from '@pkgs/cache/src/CacheProviderTypes';
|
||||
import type {CacheLookupResult} from '@pkgs/cache/src/ICacheService';
|
||||
|
||||
export function safeJsonParse<T>(value: string, logger?: CacheLogger): T | null {
|
||||
export function parseCachedValue<T>(value: string, logger?: CacheLogger): CacheLookupResult<T> {
|
||||
try {
|
||||
return JSON.parse(value);
|
||||
return {hit: true, value: JSON.parse(value)};
|
||||
} catch (error) {
|
||||
if (logger) {
|
||||
const truncatedValue = value.length > 200 ? `${value.substring(0, 200)}...` : value;
|
||||
const errorMessage = error instanceof Error ? error.message : String(error);
|
||||
logger.error({errorMessage, value: truncatedValue}, '[CacheProvider] JSON parse error');
|
||||
}
|
||||
return null;
|
||||
return {hit: false};
|
||||
}
|
||||
}
|
||||
|
||||
export function safeJsonParse<T>(value: string, logger?: CacheLogger): T | null {
|
||||
const parsed = parseCachedValue<T>(value, logger);
|
||||
return parsed.hit ? parsed.value : null;
|
||||
}
|
||||
|
||||
export function serializeValue<T>(value: T): string {
|
||||
return JSON.stringify(value);
|
||||
}
|
||||
|
||||
+33
-7
@@ -1,13 +1,19 @@
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
const CACHE_INFLIGHT_MAX_ENTRIES = 10000;
|
||||
|
||||
interface CacheMSetEntry<T> {
|
||||
key: string;
|
||||
value: T;
|
||||
ttlSeconds?: number;
|
||||
}
|
||||
|
||||
export type CacheLookupResult<T> = {hit: true; value: T} | {hit: false};
|
||||
|
||||
export abstract class ICacheService {
|
||||
abstract get<T>(key: string): Promise<T | null>;
|
||||
private readonly inflightValues = new Map<string, Promise<unknown>>();
|
||||
|
||||
abstract getEntry<T>(key: string): Promise<CacheLookupResult<T>>;
|
||||
|
||||
abstract set<T>(key: string, value: T, ttlSeconds?: number): Promise<void>;
|
||||
|
||||
@@ -45,13 +51,33 @@ export abstract class ICacheService {
|
||||
|
||||
abstract sismember(key: string, member: string): Promise<boolean>;
|
||||
|
||||
async get<T>(key: string): Promise<T | null> {
|
||||
const entry = await this.getEntry<T>(key);
|
||||
return entry.hit ? entry.value : null;
|
||||
}
|
||||
|
||||
async getOrSet<T>(key: string, valueFactory: () => Promise<T>, ttlSeconds?: number): Promise<T> {
|
||||
const existingValue = await this.get<T>(key);
|
||||
if (existingValue !== null) {
|
||||
return existingValue;
|
||||
const existing = await this.getEntry<T>(key);
|
||||
if (existing.hit) {
|
||||
return existing.value;
|
||||
}
|
||||
const newValue = await valueFactory();
|
||||
await this.set(key, newValue, ttlSeconds);
|
||||
return newValue;
|
||||
const inflight = this.inflightValues.get(key);
|
||||
if (inflight) {
|
||||
return (await inflight) as T;
|
||||
}
|
||||
if (this.inflightValues.size >= CACHE_INFLIGHT_MAX_ENTRIES) {
|
||||
return await this.produceAndStore(key, valueFactory, ttlSeconds);
|
||||
}
|
||||
const pending = this.produceAndStore(key, valueFactory, ttlSeconds).finally(() => {
|
||||
this.inflightValues.delete(key);
|
||||
});
|
||||
this.inflightValues.set(key, pending);
|
||||
return await pending;
|
||||
}
|
||||
|
||||
private async produceAndStore<T>(key: string, valueFactory: () => Promise<T>, ttlSeconds?: number): Promise<T> {
|
||||
const value = await valueFactory();
|
||||
await this.set(key, value, ttlSeconds);
|
||||
return value;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,93 @@
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
import {InMemoryProvider} from '@pkgs/cache/src/providers/InMemoryProvider';
|
||||
import {KVCacheProvider} from '@pkgs/cache/src/providers/KVCacheProvider';
|
||||
import type {IKVProvider} from '@pkgs/kv_client/src/IKVProvider';
|
||||
import {describe, expect, it, vi} from 'vitest';
|
||||
|
||||
function createKVCacheProvider(): {provider: KVCacheProvider; store: Map<string, string>} {
|
||||
const store = new Map<string, string>();
|
||||
const client = {
|
||||
get: async (key: string) => store.get(key) ?? null,
|
||||
set: async (key: string, value: string) => {
|
||||
store.set(key, value);
|
||||
return 'OK';
|
||||
},
|
||||
setex: async (key: string, _ttlSeconds: number, value: string) => {
|
||||
store.set(key, value);
|
||||
},
|
||||
} as unknown as IKVProvider;
|
||||
return {provider: new KVCacheProvider({client}), store};
|
||||
}
|
||||
|
||||
function delay(ms: number): Promise<void> {
|
||||
return new Promise((resolve) => setTimeout(resolve, ms));
|
||||
}
|
||||
|
||||
describe('ICacheService.getOrSet', () => {
|
||||
it('runs the factory once for concurrent callers on the same key', async () => {
|
||||
const cache = new InMemoryProvider();
|
||||
const factory = vi.fn(async () => {
|
||||
await delay(10);
|
||||
return 7;
|
||||
});
|
||||
const results = await Promise.all([
|
||||
cache.getOrSet('key', factory),
|
||||
cache.getOrSet('key', factory),
|
||||
cache.getOrSet('key', factory),
|
||||
]);
|
||||
expect(results).toEqual([7, 7, 7]);
|
||||
expect(factory).toHaveBeenCalledTimes(1);
|
||||
await expect(cache.get('key')).resolves.toBe(7);
|
||||
});
|
||||
|
||||
it('does not coalesce concurrent callers on different keys', async () => {
|
||||
const cache = new InMemoryProvider();
|
||||
const factory = vi.fn(async () => {
|
||||
await delay(10);
|
||||
return 1;
|
||||
});
|
||||
await Promise.all([cache.getOrSet('a', factory), cache.getOrSet('b', factory)]);
|
||||
expect(factory).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
|
||||
it('caches a null factory result and serves it as a hit', async () => {
|
||||
const cache = new InMemoryProvider();
|
||||
const factory = vi.fn(async () => null);
|
||||
await expect(cache.getOrSet<number | null>('key', factory, 60)).resolves.toBeNull();
|
||||
await expect(cache.getOrSet<number | null>('key', factory, 60)).resolves.toBeNull();
|
||||
expect(factory).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it('serves a stored json null from the kv provider as a hit', async () => {
|
||||
const {provider, store} = createKVCacheProvider();
|
||||
const factory = vi.fn(async () => null);
|
||||
await expect(provider.getOrSet<number | null>('key', factory, 60)).resolves.toBeNull();
|
||||
expect(store.get('key')).toBe('null');
|
||||
await expect(provider.getOrSet<number | null>('key', factory, 60)).resolves.toBeNull();
|
||||
expect(factory).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it('treats an unparseable stored value as a miss', async () => {
|
||||
const {provider, store} = createKVCacheProvider();
|
||||
store.set('key', '{not json');
|
||||
const factory = vi.fn(async () => 3);
|
||||
await expect(provider.getOrSet('key', factory, 60)).resolves.toBe(3);
|
||||
expect(factory).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it('rejects every waiter and retries on the next call when the factory fails', async () => {
|
||||
const cache = new InMemoryProvider();
|
||||
const failing = vi.fn(async () => {
|
||||
await delay(10);
|
||||
throw new Error('factory failed');
|
||||
});
|
||||
const settled = await Promise.allSettled([cache.getOrSet('key', failing), cache.getOrSet('key', failing)]);
|
||||
expect(settled.map((result) => result.status)).toEqual(['rejected', 'rejected']);
|
||||
expect(failing).toHaveBeenCalledTimes(1);
|
||||
await expect(cache.exists('key')).resolves.toBe(false);
|
||||
const succeeding = vi.fn(async () => 11);
|
||||
await expect(cache.getOrSet('key', succeeding)).resolves.toBe(11);
|
||||
expect(succeeding).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
});
|
||||
+5
-5
@@ -6,7 +6,7 @@ import {
|
||||
validateLockKey,
|
||||
validateLockToken,
|
||||
} from '@pkgs/cache/src/CacheLockValidation';
|
||||
import {ICacheService} from '@pkgs/cache/src/ICacheService';
|
||||
import {type CacheLookupResult, ICacheService} from '@pkgs/cache/src/ICacheService';
|
||||
|
||||
interface CacheEntry<T> {
|
||||
value: T;
|
||||
@@ -67,14 +67,14 @@ export class InMemoryProvider extends ICacheService {
|
||||
}
|
||||
}
|
||||
|
||||
async get<T>(key: string): Promise<T | null> {
|
||||
async getEntry<T>(key: string): Promise<CacheLookupResult<T>> {
|
||||
const entry = this.cache.get(key) as CacheEntry<T> | undefined;
|
||||
if (!entry) return null;
|
||||
if (!entry) return {hit: false};
|
||||
if (this.isExpired(entry)) {
|
||||
this.cache.delete(key);
|
||||
return null;
|
||||
return {hit: false};
|
||||
}
|
||||
return entry.value;
|
||||
return {hit: true, value: entry.value};
|
||||
}
|
||||
|
||||
async set<T>(key: string, value: T, ttlSeconds?: number): Promise<void> {
|
||||
|
||||
+7
-7
@@ -8,8 +8,8 @@ import {
|
||||
validateLockToken,
|
||||
} from '@pkgs/cache/src/CacheLockValidation';
|
||||
import type {CacheLogger, CacheTelemetry} from '@pkgs/cache/src/CacheProviderTypes';
|
||||
import {safeJsonParse, serializeValue} from '@pkgs/cache/src/CacheSerialization';
|
||||
import {ICacheService} from '@pkgs/cache/src/ICacheService';
|
||||
import {parseCachedValue, safeJsonParse, serializeValue} from '@pkgs/cache/src/CacheSerialization';
|
||||
import {type CacheLookupResult, ICacheService} from '@pkgs/cache/src/ICacheService';
|
||||
import type {IKVProvider} from '@pkgs/kv_client/src/IKVProvider';
|
||||
|
||||
interface KVCacheProviderConfig {
|
||||
@@ -69,16 +69,16 @@ export class KVCacheProvider extends ICacheService {
|
||||
}
|
||||
}
|
||||
|
||||
async get<T>(key: string): Promise<T | null> {
|
||||
async getEntry<T>(key: string): Promise<CacheLookupResult<T>> {
|
||||
return this.instrumented(
|
||||
'get',
|
||||
key,
|
||||
async () => {
|
||||
async (): Promise<CacheLookupResult<T>> => {
|
||||
const value = await this.client.get(key);
|
||||
if (value == null) return null;
|
||||
return safeJsonParse<T>(value, this.logger);
|
||||
if (value == null) return {hit: false};
|
||||
return parseCachedValue<T>(value, this.logger);
|
||||
},
|
||||
(result) => (result == null ? 'miss' : 'hit'),
|
||||
(result) => (result.hit ? 'hit' : 'miss'),
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
+27
@@ -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
+7
@@ -640,6 +640,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/captcha:
|
||||
dependencies:
|
||||
@@ -10772,6 +10778,7 @@ packages:
|
||||
tsconfck@3.1.6:
|
||||
resolution: {integrity: sha512-ks6Vjr/jEw0P1gmOVwutM3B7fWxoWBL2KRDb1JfqGVawBmO5UsvmWOQFGHBPl5yxYz4eERr19E6L7NMv+Fej4w==}
|
||||
engines: {node: ^18 || >=20}
|
||||
deprecated: unmaintained
|
||||
hasBin: true
|
||||
peerDependencies:
|
||||
typescript: ^5.0.0
|
||||
|
||||
Reference in New Issue
Block a user