From 2d6b01691c0e33ce2ce3ae2ef56ee16f41891a08 Mon Sep 17 00:00:00 2001 From: Ben Date: Sun, 7 Jun 2026 21:32:03 -0400 Subject: [PATCH] Cache hot API responses in Redis --- src/app.controller.ts | 41 +++++++++--- src/services/redis-messaging.service.spec.ts | 9 ++- src/services/redis-messaging.service.ts | 67 +++++++++++++++++--- 3 files changed, 96 insertions(+), 21 deletions(-) diff --git a/src/app.controller.ts b/src/app.controller.ts index 71498e3..173b842 100644 --- a/src/app.controller.ts +++ b/src/app.controller.ts @@ -12,6 +12,7 @@ import { BitcoinRpcService } from './services/bitcoin-rpc.service'; import { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-report.view'; import { StratumV2Service } from './services/stratum-v2.service'; import { ShareAccountingService } from './ORM/share-accounting/share-accounting.service'; +import { RedisMessagingService } from './services/redis-messaging.service'; @Controller() export class AppController { @@ -27,7 +28,8 @@ export class AppController { private readonly addressSettingsService: AddressSettingsService, private readonly userAgentReportService: UserAgentReportService, private readonly stratumV2Service: StratumV2Service, - private readonly shareAccountingService: ShareAccountingService + private readonly shareAccountingService: ShareAccountingService, + private readonly redisMessagingService: RedisMessagingService ) { } @Get('info') @@ -35,7 +37,7 @@ export class AppController { const CACHE_KEY = 'SITE_INFO'; - const cachedResult = await this.cacheManager.get(CACHE_KEY); + const cachedResult = await this.getCached(CACHE_KEY); if (cachedResult != null) { return cachedResult; @@ -86,7 +88,7 @@ export class AppController { }; // Keep online miner counts responsive after reconnect cleanup. - await this.cacheManager.set(CACHE_KEY, data, 15 * 1000); + await this.setCached(CACHE_KEY, data, 15 * 1000); return data; @@ -95,7 +97,7 @@ export class AppController { @Get('info/accounting') public async infoAccounting() { const CACHE_KEY = 'SITE_ACCOUNTING'; - const cachedResult = await this.cacheManager.get(CACHE_KEY); + const cachedResult = await this.getCached(CACHE_KEY); if (cachedResult != null) { return cachedResult; @@ -104,7 +106,7 @@ export class AppController { const data = await this.shareAccountingService.getPoolSummary(); //15 sec - await this.cacheManager.set(CACHE_KEY, data, 15 * 1000); + await this.setCached(CACHE_KEY, data, 15 * 1000); return data; } @@ -113,7 +115,7 @@ export class AppController { public async pool() { const CACHE_KEY = 'POOL_INFO'; - const cachedResult = await this.cacheManager.get(CACHE_KEY); + const cachedResult = await this.getCached(CACHE_KEY); if (cachedResult != null) { return cachedResult; @@ -135,7 +137,7 @@ export class AppController { } // Keep online miner counts responsive after reconnect cleanup. - await this.cacheManager.set(CACHE_KEY, data, 15 * 1000); + await this.setCached(CACHE_KEY, data, 15 * 1000); return data; } @@ -150,7 +152,7 @@ export class AppController { const CACHE_KEY = 'SITE_HASHRATE_GRAPH'; - const cachedResult = await this.cacheManager.get(CACHE_KEY); + const cachedResult = await this.getCached(CACHE_KEY); if (cachedResult != null) { return cachedResult; @@ -159,11 +161,32 @@ export class AppController { const chartData = await this.clientStatisticsService.getChartDataForSite(); //10 min - await this.cacheManager.set(CACHE_KEY, chartData, 10 * 60 * 1000); + await this.setCached(CACHE_KEY, chartData, 10 * 60 * 1000); return chartData; } + private async getCached(key: string): Promise { + const shared = await this.redisMessagingService.getJsonCache(`api:${key}`).catch(error => { + console.error(`Shared API cache read failed for ${key}: ${error.message}`); + return null; + }); + if (shared != null) { + return shared; + } + + return await this.cacheManager.get(key) ?? null; + } + + private async setCached(key: string, value: unknown, ttlMs: number): Promise { + await Promise.all([ + this.cacheManager.set(key, value, ttlMs), + this.redisMessagingService.setJsonCache(`api:${key}`, value, ttlMs).catch(error => { + console.error(`Shared API cache write failed for ${key}: ${error.message}`); + }), + ]); + } + } diff --git a/src/services/redis-messaging.service.spec.ts b/src/services/redis-messaging.service.spec.ts index b35366f..a7cc99a 100644 --- a/src/services/redis-messaging.service.spec.ts +++ b/src/services/redis-messaging.service.spec.ts @@ -133,6 +133,7 @@ function createRedisClient() { return Promise.resolve('OK'); }), get: jest.fn((key: string) => Promise.resolve(store.get(key) ?? null)), + mGet: jest.fn((keys: string[]) => Promise.resolve(keys.map(key => store.get(key) ?? null))), del: jest.fn((...args: (string | string[])[]) => { const keys = args.flatMap(key => Array.isArray(key) ? key : [key]); let deleted = 0; @@ -148,9 +149,13 @@ function createRedisClient() { sets.set(key, set); return Promise.resolve(1); }), - sRem: jest.fn((key: string, value: string) => { + sRem: jest.fn((key: string, value: string | string[]) => { const set = sets.get(key); - const deleted = set?.delete(value) ? 1 : 0; + const values = Array.isArray(value) ? value : [value]; + let deleted = 0; + values.forEach(item => { + deleted += set?.delete(item) ? 1 : 0; + }); return Promise.resolve(deleted); }), sMembers: jest.fn((key: string) => Promise.resolve([...sets.get(key) ?? []])), diff --git a/src/services/redis-messaging.service.ts b/src/services/redis-messaging.service.ts index 25d5a55..123e55b 100644 --- a/src/services/redis-messaging.service.ts +++ b/src/services/redis-messaging.service.ts @@ -13,6 +13,7 @@ const blockTemplateKey = (height: number) => `block-template:${height}`; const CLIENT_PRESENCE_ALL_KEY = 'client-presence:all'; const clientPresenceKey = (clientId: string) => `client-presence:${clientId}`; const clientPresenceAddressKey = (address: string) => `client-presence:address:${address}`; +const jsonCacheKey = (key: string) => `json-cache:${key}`; export interface ClientPresence { clientId: string; @@ -218,6 +219,37 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy { await this.deleteKeys(batch); } + public async getJsonCache(key: string): Promise { + if (!await this.ensureConnected()) { + return null; + } + + const value = await this.publisher.get(jsonCacheKey(key)); + if (value == null) { + return null; + } + + try { + return JSON.parse(value as string) as T; + } catch (error) { + console.error(`Invalid Redis JSON cache for ${key}: ${error.message}`); + await this.publisher.del(jsonCacheKey(key)); + return null; + } + } + + public async setJsonCache(key: string, value: unknown, ttlMs: number): Promise { + if (!await this.ensureConnected() || ttlMs <= 0) { + return; + } + + await this.publisher.setEx( + jsonCacheKey(key), + Math.max(1, Math.ceil(ttlMs / 1000)), + JSON.stringify(value), + ); + } + private async ensureConnected(): Promise { if (!this.connected) { try { @@ -236,18 +268,33 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy { return []; } - const presences = await Promise.all(clientIds.map(async clientId => { - const presence = await this.getClientPresence(clientId); - if (presence == null) { - await this.publisher.sRem(setKey, clientId); - if (setKey !== CLIENT_PRESENCE_ALL_KEY) { - await this.publisher.sRem(CLIENT_PRESENCE_ALL_KEY, clientId); + const presences: ClientPresence[] = []; + const staleClientIds: string[] = []; + for (let i = 0; i < clientIds.length; i += 1000) { + const chunk = clientIds.slice(i, i + 1000); + const values = await this.publisher.mGet(chunk.map(clientPresenceKey)); + values.forEach((value, index) => { + const presence = this.parseClientPresence(value); + if (presence == null) { + staleClientIds.push(chunk[index]); + return; } - } - return presence; - })); + presences.push(presence); + }); + } - return presences.filter((presence): presence is ClientPresence => presence != null); + for (let i = 0; i < staleClientIds.length; i += 1000) { + const staleChunk = staleClientIds.slice(i, i + 1000); + if (staleChunk.length === 0) { + continue; + } + await this.publisher.sRem(setKey, staleChunk); + if (setKey !== CLIENT_PRESENCE_ALL_KEY) { + await this.publisher.sRem(CLIENT_PRESENCE_ALL_KEY, staleChunk); + } + } + + return presences; } private parseClientPresence(value: unknown): ClientPresence | null {