diff --git a/fluxer_api/pkgs/cache/package.json b/fluxer_api/pkgs/cache/package.json index 34ec62bc9..60bb299b0 100644 --- a/fluxer_api/pkgs/cache/package.json +++ b/fluxer_api/pkgs/cache/package.json @@ -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:" } } diff --git a/fluxer_api/pkgs/cache/src/CacheSerialization.ts b/fluxer_api/pkgs/cache/src/CacheSerialization.ts index a3bfacba1..a8102df42 100644 --- a/fluxer_api/pkgs/cache/src/CacheSerialization.ts +++ b/fluxer_api/pkgs/cache/src/CacheSerialization.ts @@ -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(value: string, logger?: CacheLogger): T | null { +export function parseCachedValue(value: string, logger?: CacheLogger): CacheLookupResult { 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(value: string, logger?: CacheLogger): T | null { + const parsed = parseCachedValue(value, logger); + return parsed.hit ? parsed.value : null; +} + export function serializeValue(value: T): string { return JSON.stringify(value); } diff --git a/fluxer_api/pkgs/cache/src/ICacheService.ts b/fluxer_api/pkgs/cache/src/ICacheService.ts index afd2c2820..e34e6fadb 100644 --- a/fluxer_api/pkgs/cache/src/ICacheService.ts +++ b/fluxer_api/pkgs/cache/src/ICacheService.ts @@ -1,13 +1,19 @@ // SPDX-License-Identifier: AGPL-3.0-or-later +const CACHE_INFLIGHT_MAX_ENTRIES = 10000; + interface CacheMSetEntry { key: string; value: T; ttlSeconds?: number; } +export type CacheLookupResult = {hit: true; value: T} | {hit: false}; + export abstract class ICacheService { - abstract get(key: string): Promise; + private readonly inflightValues = new Map>(); + + abstract getEntry(key: string): Promise>; abstract set(key: string, value: T, ttlSeconds?: number): Promise; @@ -45,13 +51,33 @@ export abstract class ICacheService { abstract sismember(key: string, member: string): Promise; + async get(key: string): Promise { + const entry = await this.getEntry(key); + return entry.hit ? entry.value : null; + } + async getOrSet(key: string, valueFactory: () => Promise, ttlSeconds?: number): Promise { - const existingValue = await this.get(key); - if (existingValue !== null) { - return existingValue; + const existing = await this.getEntry(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(key: string, valueFactory: () => Promise, ttlSeconds?: number): Promise { + const value = await valueFactory(); + await this.set(key, value, ttlSeconds); + return value; } } diff --git a/fluxer_api/pkgs/cache/src/__tests__/CacheGetOrSet.test.ts b/fluxer_api/pkgs/cache/src/__tests__/CacheGetOrSet.test.ts new file mode 100644 index 000000000..73d511af4 --- /dev/null +++ b/fluxer_api/pkgs/cache/src/__tests__/CacheGetOrSet.test.ts @@ -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} { + const store = new Map(); + 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 { + 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('key', factory, 60)).resolves.toBeNull(); + await expect(cache.getOrSet('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('key', factory, 60)).resolves.toBeNull(); + expect(store.get('key')).toBe('null'); + await expect(provider.getOrSet('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); + }); +}); diff --git a/fluxer_api/pkgs/cache/src/providers/InMemoryProvider.ts b/fluxer_api/pkgs/cache/src/providers/InMemoryProvider.ts index 132a25420..0e7775820 100644 --- a/fluxer_api/pkgs/cache/src/providers/InMemoryProvider.ts +++ b/fluxer_api/pkgs/cache/src/providers/InMemoryProvider.ts @@ -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 { value: T; @@ -67,14 +67,14 @@ export class InMemoryProvider extends ICacheService { } } - async get(key: string): Promise { + async getEntry(key: string): Promise> { const entry = this.cache.get(key) as CacheEntry | 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(key: string, value: T, ttlSeconds?: number): Promise { diff --git a/fluxer_api/pkgs/cache/src/providers/KVCacheProvider.ts b/fluxer_api/pkgs/cache/src/providers/KVCacheProvider.ts index f5bbd7cc2..a388f849f 100644 --- a/fluxer_api/pkgs/cache/src/providers/KVCacheProvider.ts +++ b/fluxer_api/pkgs/cache/src/providers/KVCacheProvider.ts @@ -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(key: string): Promise { + async getEntry(key: string): Promise> { return this.instrumented( 'get', key, - async () => { + async (): Promise> => { const value = await this.client.get(key); - if (value == null) return null; - return safeJsonParse(value, this.logger); + if (value == null) return {hit: false}; + return parseCachedValue(value, this.logger); }, - (result) => (result == null ? 'miss' : 'hit'), + (result) => (result.hit ? 'hit' : 'miss'), ); } diff --git a/fluxer_api/pkgs/cache/vitest.config.ts b/fluxer_api/pkgs/cache/vitest.config.ts new file mode 100644 index 000000000..08796abc6 --- /dev/null +++ b/fluxer_api/pkgs/cache/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 d1e89db1d..343e6ea1c 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -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