diff --git a/.env.example b/.env.example index 7614a73..27412f7 100644 --- a/.env.example +++ b/.env.example @@ -170,6 +170,10 @@ SHARE_ACCOUNTING_FLUSH_INTERVAL_MS=25 SHARE_ACCOUNTING_MAX_QUEUE_SIZE=50000 SHARE_ACCOUNTING_SUMMARY_CACHE_MS=2500 SHARE_ACCOUNTING_SUMMARY_CACHE_MAX=10000 +SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS=300000 +SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS=3600000 +SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS=30000 +SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS=5000 CLIENT_HASHRATE_PERSIST_INTERVAL_MS=60000 SHARE_ROLLUP_ENABLED=true SHARE_ROLLUP_INTERVAL_MS=60000 diff --git a/docker-compose.external-db.yml b/docker-compose.external-db.yml index 88d8506..60e686f 100644 --- a/docker-compose.external-db.yml +++ b/docker-compose.external-db.yml @@ -106,6 +106,10 @@ services: SHARE_ACCOUNTING_MAX_QUEUE_SIZE: ${SHARE_ACCOUNTING_MAX_QUEUE_SIZE:-50000} SHARE_ACCOUNTING_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MS:-2500} SHARE_ACCOUNTING_SUMMARY_CACHE_MAX: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MAX:-10000} + SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS:-300000} + SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS:-3600000} + SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS:-30000} + SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS:-5000} SHARE_ROLLUP_ENABLED: ${SHARE_ROLLUP_ENABLED:-true} SHARE_ROLLUP_INTERVAL_MS: ${SHARE_ROLLUP_INTERVAL_MS:-60000} SHARE_ROLLUP_SAFETY_LAG_SECONDS: ${SHARE_ROLLUP_SAFETY_LAG_SECONDS:-30} diff --git a/docker-compose.yml b/docker-compose.yml index 48697e6..14ef2b0 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -110,6 +110,10 @@ services: API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000} API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100} API_MAX_CONNECTIONS_PER_WORKER: ${API_MAX_CONNECTIONS_PER_WORKER:-500} + SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS:-300000} + SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS:-3600000} + SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS:-30000} + SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS:-5000} PM2_ENABLED: ${PM2_ENABLED:-true} STRATUM_WORKERS: ${STRATUM_WORKERS:-auto} STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330} diff --git a/full-setup/docker-compose-mainnet.yml b/full-setup/docker-compose-mainnet.yml index 279aa95..10559d0 100644 --- a/full-setup/docker-compose-mainnet.yml +++ b/full-setup/docker-compose-mainnet.yml @@ -114,6 +114,10 @@ services: API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000} API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100} API_MAX_CONNECTIONS_PER_WORKER: ${API_MAX_CONNECTIONS_PER_WORKER:-500} + SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS:-300000} + SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS:-3600000} + SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS:-30000} + SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS:-5000} PM2_ENABLED: "true" STRATUM_WORKERS: ${STRATUM_WORKERS:-2} STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1} diff --git a/src/ORM/share-accounting/share-accounting.service.spec.ts b/src/ORM/share-accounting/share-accounting.service.spec.ts index 6d6bbaa..18e40b1 100644 --- a/src/ORM/share-accounting/share-accounting.service.spec.ts +++ b/src/ORM/share-accounting/share-accounting.service.spec.ts @@ -365,11 +365,65 @@ describe('ShareAccountingService', () => { expect(repository.query).toHaveBeenCalledTimes(1); expect(redis.setJsonCache).toHaveBeenCalledWith( expect.stringContaining('accounting:summary:'), - expect.objectContaining({ totalAcceptedShares: 1 }), - 30000, + expect.objectContaining({ + schemaVersion: 1, + refreshedAtMs: expect.any(Number), + value: expect.objectContaining({ totalAcceptedShares: 1 }), + }), + 3600000, ); }); + it('serves a stale shared summary while one API worker refreshes it', async () => { + process.env.SHARE_ACCOUNTING_SUMMARY_CACHE_MS = '0'; + process.env.SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS = '100'; + const staleSummary = { + ...new ShareAccountingService({} as any).emptySummary(), + totalAcceptedShares: 7, + }; + const repository = { + query: jest.fn().mockResolvedValueOnce([{ + totalAcceptedShares: '8', + totalCreditedDifficulty: '256', + acceptedSharesLast10Minutes: '1', + creditedDifficultyLast10Minutes: '32', + acceptedSharesLastHour: '8', + creditedDifficultyLastHour: '256', + acceptedSharesLastDay: '8', + creditedDifficultyLastDay: '256', + hashRateLast10Minutes: '1', + hashRateLastHour: '1', + latestShareAt: null, + }]), + }; + const redis = { + getJsonCache: jest.fn().mockResolvedValue({ + schemaVersion: 1, + refreshedAtMs: Date.now() - 1000, + value: staleSummary, + }), + setJsonCache: jest.fn().mockResolvedValue(undefined), + tryAcquireJsonCacheLock: jest.fn().mockResolvedValue(true), + releaseJsonCacheLock: jest.fn().mockResolvedValue(undefined), + }; + const service = new ShareAccountingService(repository as any, redis as any); + + await expect(service.getAddressSummary('bc1qstale')).resolves.toEqual(staleSummary); + await new Promise(resolve => setImmediate(resolve)); + + expect(repository.query).toHaveBeenCalledTimes(1); + expect(redis.tryAcquireJsonCacheLock).toHaveBeenCalledTimes(1); + expect(redis.setJsonCache).toHaveBeenCalledWith( + expect.stringContaining('accounting:summary:'), + expect.objectContaining({ + schemaVersion: 1, + value: expect.objectContaining({ totalAcceptedShares: 8 }), + }), + 3600000, + ); + expect(redis.releaseJsonCacheLock).toHaveBeenCalledTimes(1); + }); + it('should refresh pool summaries from completed rollup buckets and current round rollups', async () => { const redis = { getJsonCache: jest.fn().mockResolvedValue(null), diff --git a/src/ORM/share-accounting/share-accounting.service.ts b/src/ORM/share-accounting/share-accounting.service.ts index b1401ce..284cece 100644 --- a/src/ORM/share-accounting/share-accounting.service.ts +++ b/src/ORM/share-accounting/share-accounting.service.ts @@ -1,5 +1,6 @@ import { Injectable, OnModuleDestroy, OnModuleInit } from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; +import { randomUUID } from 'crypto'; import { Repository } from 'typeorm'; import { AcceptedShareEntity } from '../accepted-share/accepted-share.entity'; @@ -88,7 +89,10 @@ const DEFAULT_FLUSH_INTERVAL_MS = 25; const DEFAULT_MAX_QUEUE_SIZE = 50000; const DEFAULT_SUMMARY_CACHE_MS = 2500; const DEFAULT_SUMMARY_CACHE_MAX = 10000; -const DEFAULT_REDIS_SUMMARY_CACHE_MS = 30000; +const DEFAULT_REDIS_SUMMARY_CACHE_MS = 5 * 60 * 1000; +const DEFAULT_REDIS_SUMMARY_STALE_CACHE_MS = 60 * 60 * 1000; +const DEFAULT_REDIS_SUMMARY_LOCK_MS = 30 * 1000; +const DEFAULT_REDIS_SUMMARY_LOCK_WAIT_MS = 5 * 1000; const DEFAULT_ROLLUP_INTERVAL_MS = 60000; const DEFAULT_ROLLUP_SAFETY_LAG_SECONDS = 10; const DEFAULT_ROLLUP_MAX_SHARES_PER_BATCH = 5000000; @@ -103,6 +107,7 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy { private activeRollup: Promise | null = null; private summaryCache = new Map(); private summaryInFlight = new Map>(); + private redisSummaryRefreshInFlight = new Map>(); private readonly poolSummaryCacheKey = 'accounting:pool-summary'; private readonly batchSize = this.readPositiveInt('SHARE_ACCOUNTING_BATCH_SIZE', DEFAULT_BATCH_SIZE); private readonly flushIntervalMs = this.readPositiveInt('SHARE_ACCOUNTING_FLUSH_INTERVAL_MS', DEFAULT_FLUSH_INTERVAL_MS); @@ -110,6 +115,12 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy { private readonly summaryCacheMs = this.readNonNegativeInt('SHARE_ACCOUNTING_SUMMARY_CACHE_MS', DEFAULT_SUMMARY_CACHE_MS); private readonly summaryCacheMax = this.readPositiveInt('SHARE_ACCOUNTING_SUMMARY_CACHE_MAX', DEFAULT_SUMMARY_CACHE_MAX); private readonly redisSummaryCacheMs = this.readNonNegativeInt('SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS', DEFAULT_REDIS_SUMMARY_CACHE_MS); + private readonly redisSummaryStaleCacheMs = Math.max( + this.redisSummaryCacheMs, + this.readPositiveInt('SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS', DEFAULT_REDIS_SUMMARY_STALE_CACHE_MS), + ); + private readonly redisSummaryLockMs = this.readPositiveInt('SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS', DEFAULT_REDIS_SUMMARY_LOCK_MS); + private readonly redisSummaryLockWaitMs = this.readPositiveInt('SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS', DEFAULT_REDIS_SUMMARY_LOCK_WAIT_MS); private readonly shareRollupEnabled = this.readBoolean('SHARE_ROLLUP_ENABLED', true); private readonly shareRollupIntervalMs = this.readPositiveInt('SHARE_ROLLUP_INTERVAL_MS', DEFAULT_ROLLUP_INTERVAL_MS); private readonly shareRollupSafetyLagSeconds = this.readNonNegativeInt('SHARE_ROLLUP_SAFETY_LAG_SECONDS', DEFAULT_ROLLUP_SAFETY_LAG_SECONDS); @@ -597,25 +608,175 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy { private async loadCachedSummary(filter: AccountingFilter, cacheKey: string): Promise { const redisCacheKey = `accounting:summary:${cacheKey}`; + if (this.redisMessagingService == null || this.redisSummaryCacheMs <= 0) { + return this.loadSummary(filter); + } + + const cached = await this.readRedisSummaryCache(redisCacheKey); + if (cached != null) { + if (Date.now() - cached.refreshedAtMs >= this.redisSummaryCacheMs) { + this.refreshRedisSummaryCacheInBackground(filter, redisCacheKey); + } + return cached.value; + } + + return this.loadColdRedisSummaryCache(filter, redisCacheKey); + } + + private async readRedisSummaryCache(redisCacheKey: string): Promise { const cached = await this.redisMessagingService - ?.getJsonCache(redisCacheKey) + ?.getJsonCache(redisCacheKey) .catch(error => { console.error(`Share accounting summary cache read failed: ${error.message}`); return null; }); - if (cached != null) { + if (cached == null) { + return null; + } + + if (this.isRedisSummaryCacheEntry(cached)) { return cached; } + // Cache entries written by older workers have no timestamp. Their old + // short Redis TTL still bounds how long they can be treated as fresh. + return { + schemaVersion: 1, + refreshedAtMs: Date.now(), + value: cached, + }; + } + + private async loadColdRedisSummaryCache( + filter: AccountingFilter, + redisCacheKey: string, + ): Promise { + if (!this.supportsRedisSummaryLock()) { + const summary = await this.loadSummary(filter); + await this.storeRedisSummaryCache(redisCacheKey, summary); + return summary; + } + + const owner = randomUUID(); + const acquired = await this.tryAcquireRedisSummaryLock(redisCacheKey, owner); + if (acquired === true) { + return this.loadAndStoreRedisSummary(filter, redisCacheKey, owner); + } + if (acquired == null) { + const summary = await this.loadSummary(filter); + await this.storeRedisSummaryCache(redisCacheKey, summary); + return summary; + } + + const deadline = Date.now() + this.redisSummaryLockWaitMs; + while (Date.now() < deadline) { + await new Promise(resolve => setTimeout(resolve, 50)); + const cached = await this.readRedisSummaryCache(redisCacheKey); + if (cached != null) { + return cached.value; + } + } + + // Redis locking is an optimization, not an availability dependency. + // If a lock holder died or a refresh exceeded its budget, fail open. const summary = await this.loadSummary(filter); + await this.storeRedisSummaryCache(redisCacheKey, summary); + return summary; + } + + private refreshRedisSummaryCacheInBackground(filter: AccountingFilter, redisCacheKey: string): void { + if (this.redisSummaryRefreshInFlight.has(redisCacheKey)) { + return; + } + + const refresh = this.refreshRedisSummaryCache(filter, redisCacheKey) + .catch(error => { + console.error(`Share accounting summary background refresh failed: ${error.message}`); + }) + .finally(() => { + this.redisSummaryRefreshInFlight.delete(redisCacheKey); + }); + this.redisSummaryRefreshInFlight.set(redisCacheKey, refresh); + } + + private async refreshRedisSummaryCache(filter: AccountingFilter, redisCacheKey: string): Promise { + if (!this.supportsRedisSummaryLock()) { + const summary = await this.loadSummary(filter); + await this.storeRedisSummaryCache(redisCacheKey, summary); + return; + } + + const owner = randomUUID(); + if (await this.tryAcquireRedisSummaryLock(redisCacheKey, owner) !== true) { + return; + } + + await this.loadAndStoreRedisSummary(filter, redisCacheKey, owner); + } + + private async loadAndStoreRedisSummary( + filter: AccountingFilter, + redisCacheKey: string, + owner: string, + ): Promise { + try { + const summary = await this.loadSummary(filter); + await this.storeRedisSummaryCache(redisCacheKey, summary); + return summary; + } finally { + if (this.supportsRedisSummaryLock()) { + await this.redisMessagingService + .releaseJsonCacheLock(redisCacheKey, owner) + .catch(error => { + console.error(`Share accounting summary cache lock release failed: ${error.message}`); + }); + } + } + } + + private async storeRedisSummaryCache( + redisCacheKey: string, + summary: ShareAccountingSummary, + ): Promise { + const entry: RedisSummaryCacheEntry = { + schemaVersion: 1, + refreshedAtMs: Date.now(), + value: summary, + }; await this.redisMessagingService - ?.setJsonCache(redisCacheKey, summary, this.redisSummaryCacheMs) + ?.setJsonCache(redisCacheKey, entry, this.redisSummaryStaleCacheMs) .catch(error => { console.error(`Share accounting summary cache write failed: ${error.message}`); }); + } - return summary; + private async tryAcquireRedisSummaryLock(redisCacheKey: string, owner: string): Promise { + if (!this.supportsRedisSummaryLock()) { + return null; + } + + try { + return await this.redisMessagingService + .tryAcquireJsonCacheLock(redisCacheKey, owner, this.redisSummaryLockMs); + } catch (error) { + console.error(`Share accounting summary cache lock failed: ${error.message}`); + return null; + } + } + + private supportsRedisSummaryLock(): boolean { + return typeof this.redisMessagingService?.tryAcquireJsonCacheLock === 'function' + && typeof this.redisMessagingService?.releaseJsonCacheLock === 'function'; + } + + private isRedisSummaryCacheEntry( + cached: RedisSummaryCacheEntry | ShareAccountingSummary, + ): cached is RedisSummaryCacheEntry { + const candidate = cached as Partial; + return candidate.schemaVersion === 1 + && Number.isFinite(candidate.refreshedAtMs) + && candidate.value != null; } private async loadSummary(filter: AccountingFilter): Promise { @@ -924,3 +1085,9 @@ interface SummaryCacheEntry { expiresAt: number; value: ShareAccountingSummary; } + +interface RedisSummaryCacheEntry { + schemaVersion: 1; + refreshedAtMs: number; + value: ShareAccountingSummary; +} diff --git a/src/services/redis-messaging.service.spec.ts b/src/services/redis-messaging.service.spec.ts index c51ac28..166cea8 100644 --- a/src/services/redis-messaging.service.spec.ts +++ b/src/services/redis-messaging.service.spec.ts @@ -120,6 +120,19 @@ describe('RedisMessagingService', () => { consoleSpy.mockRestore(); }); + it('uses an owner token when locking JSON cache refreshes', async () => { + await service.connect(); + + await expect(service.tryAcquireJsonCacheLock('summary', 'owner-a', 1000)).resolves.toBe(true); + await expect(service.tryAcquireJsonCacheLock('summary', 'owner-b', 1000)).resolves.toBe(false); + + await service.releaseJsonCacheLock('summary', 'owner-b'); + await expect(service.tryAcquireJsonCacheLock('summary', 'owner-b', 1000)).resolves.toBe(false); + + await service.releaseJsonCacheLock('summary', 'owner-a'); + await expect(service.tryAcquireJsonCacheLock('summary', 'owner-b', 1000)).resolves.toBe(true); + }); + it('stores payout variants and same-height reorg templates under distinct tip keys', async () => { await service.connect(); const firstHash = '11'.repeat(32); @@ -493,7 +506,10 @@ function createRedisClient() { connect: jest.fn().mockResolvedValue(undefined), quit: jest.fn().mockResolvedValue(undefined), on: jest.fn(), - set: jest.fn((key: string, value: string) => { + set: jest.fn((key: string, value: string, options?: { NX?: boolean }) => { + if (options?.NX && store.has(key)) { + return Promise.resolve(null); + } store.set(key, value); return Promise.resolve('OK'); }), @@ -512,6 +528,15 @@ function createRedisClient() { }); return Promise.resolve(deleted); }), + eval: jest.fn((_script: string, options: { keys: string[], arguments: string[] }) => { + const [key] = options.keys; + const [owner] = options.arguments; + if (store.get(key) !== owner) { + return Promise.resolve(0); + } + store.delete(key); + return Promise.resolve(1); + }), sAdd: jest.fn((key: string, value: string) => { const set = sets.get(key) ?? new Set(); set.add(value); diff --git a/src/services/redis-messaging.service.ts b/src/services/redis-messaging.service.ts index 659e0b6..e97f8c7 100644 --- a/src/services/redis-messaging.service.ts +++ b/src/services/redis-messaging.service.ts @@ -32,6 +32,7 @@ const blockTemplateHeightPointerKey = (height: number, payoutMode: PayoutMode) = const blockTemplateLatestPointerKey = (payoutMode: PayoutMode) => `${BLOCK_TEMPLATE_LATEST_KEY}:${payoutMode}`; const jsonCacheKey = (key: string) => `json-cache:${key}`; +const jsonCacheLockKey = (key: string) => `json-cache-lock:${key}`; export interface BlockTemplateUpdate { schemaVersion: 1; @@ -524,6 +525,41 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy { ); } + public async tryAcquireJsonCacheLock(key: string, owner: string, ttlMs: number): Promise { + if (!await this.ensureConnected()) { + return null; + } + if (owner.length === 0 || ttlMs <= 0) { + return false; + } + + const result = await this.publisher.set( + jsonCacheLockKey(key), + owner, + { NX: true, PX: Math.max(1, Math.ceil(ttlMs)) }, + ); + return result === 'OK'; + } + + public async releaseJsonCacheLock(key: string, owner: string): Promise { + if (!await this.ensureConnected() || owner.length === 0) { + return; + } + + await this.publisher.eval( + ` + if redis.call('GET', KEYS[1]) == ARGV[1] then + return redis.call('DEL', KEYS[1]) + end + return 0 + `, + { + keys: [jsonCacheLockKey(key)], + arguments: [owner], + }, + ); + } + private async readBlockTemplatePointer(value: string | null): Promise { if (value == null) { return null;