diff --git a/src/ORM/share-accounting/share-accounting.service.spec.ts b/src/ORM/share-accounting/share-accounting.service.spec.ts index f5d65cd..f0c6d6a 100644 --- a/src/ORM/share-accounting/share-accounting.service.spec.ts +++ b/src/ORM/share-accounting/share-accounting.service.spec.ts @@ -165,8 +165,89 @@ describe('ShareAccountingService', () => { expect(repository.query).toHaveBeenCalledTimes(1); }); + it('should serve API-only address accounting from the share rollup instead of empty summaries', async () => { + process.env.API_ONLY = 'true'; + const repository = { + query: jest.fn() + .mockResolvedValueOnce([{ + totalAcceptedShares: '12', + totalCreditedDifficulty: '384', + acceptedSharesLast10Minutes: '4', + creditedDifficultyLast10Minutes: '128', + acceptedSharesLastHour: '10', + creditedDifficultyLastHour: '320', + acceptedSharesLastDay: '12', + creditedDifficultyLastDay: '384', + hashRateLast10Minutes: '916259689.8', + hashRateLastHour: '381774870.2', + latestShareAt: new Date('2026-06-07T12:30:00Z'), + }]), + }; + const service = new ShareAccountingService(repository as any); + + await expect(service.getAddressSummary('bc1qapi')).resolves.toEqual(expect.objectContaining({ + totalAcceptedShares: 12, + totalCreditedDifficulty: 384, + acceptedSharesLast10Minutes: 4, + creditedDifficultyLast10Minutes: 128, + hashRateLast10Minutes: 916259689.8, + latestShareAt: '2026-06-07T12:30:00.000Z', + })); + expect(repository.query).toHaveBeenCalledWith( + expect.stringContaining('FROM "accepted_share_10m"'), + ['bc1qapi'], + ); + expect(repository.query).not.toHaveBeenCalledWith( + expect.stringContaining('FROM "accepted_share_entity"'), + expect.anything(), + ); + }); + + it('should use Redis cache for share accounting summaries across API workers', async () => { + process.env.SHARE_ACCOUNTING_SUMMARY_CACHE_MS = '0'; + const repository = { + query: jest.fn() + .mockResolvedValueOnce([{ + totalAcceptedShares: '1', + totalCreditedDifficulty: '32', + acceptedSharesLast10Minutes: '1', + creditedDifficultyLast10Minutes: '32', + acceptedSharesLastHour: '1', + creditedDifficultyLastHour: '32', + acceptedSharesLastDay: '1', + creditedDifficultyLastDay: '32', + hashRateLast10Minutes: '1', + hashRateLastHour: '1', + latestShareAt: null, + }]), + }; + const redis = { + getJsonCache: jest.fn() + .mockResolvedValueOnce(null) + .mockResolvedValueOnce({ + ...new ShareAccountingService(repository as any).emptySummary(), + totalAcceptedShares: 1, + }), + setJsonCache: jest.fn().mockResolvedValue(undefined), + }; + const service = new ShareAccountingService(repository as any, redis as any); + + await service.getAddressSummary('bc1qcached'); + await expect(service.getAddressSummary('bc1qcached')).resolves.toEqual(expect.objectContaining({ + totalAcceptedShares: 1, + })); + + expect(repository.query).toHaveBeenCalledTimes(1); + expect(redis.setJsonCache).toHaveBeenCalledWith( + expect.stringContaining('accounting:summary:'), + expect.objectContaining({ totalAcceptedShares: 1 }), + 30000, + ); + }); + it('should overlay live pool data and best share from the current round', async () => { const redis = { + getJsonCache: jest.fn().mockResolvedValue(null), setJsonCache: jest.fn().mockResolvedValue(undefined), }; const repository = { diff --git a/src/ORM/share-accounting/share-accounting.service.ts b/src/ORM/share-accounting/share-accounting.service.ts index 07a0a9b..fc964f1 100644 --- a/src/ORM/share-accounting/share-accounting.service.ts +++ b/src/ORM/share-accounting/share-accounting.service.ts @@ -69,6 +69,7 @@ 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; @Injectable() export class ShareAccountingService implements OnModuleDestroy { @@ -82,6 +83,7 @@ export class ShareAccountingService implements OnModuleDestroy { private readonly maxQueueSize = this.readPositiveInt('SHARE_ACCOUNTING_MAX_QUEUE_SIZE', DEFAULT_MAX_QUEUE_SIZE); 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); constructor( @InjectRepository(AcceptedShareEntity) @@ -269,10 +271,6 @@ export class ShareAccountingService implements OnModuleDestroy { } private async getSummary(filter: AccountingFilter): Promise { - if (process.env.API_ONLY === 'true') { - return this.emptySummary(); - } - const cacheKey = this.getSummaryCacheKey(filter); const cached = this.summaryCache.get(cacheKey); const now = Date.now(); @@ -281,7 +279,7 @@ export class ShareAccountingService implements OnModuleDestroy { return cached.value; } - const value = this.loadSummary(filter).catch(error => { + const value = this.loadCachedSummary(filter, cacheKey).catch(error => { this.summaryCache.delete(cacheKey); throw error; }); @@ -297,6 +295,29 @@ export class ShareAccountingService implements OnModuleDestroy { return value; } + private async loadCachedSummary(filter: AccountingFilter, cacheKey: string): Promise { + const redisCacheKey = `accounting:summary:${cacheKey}`; + const cached = await this.redisMessagingService + ?.getJsonCache(redisCacheKey) + .catch(error => { + console.error(`Share accounting summary cache read failed: ${error.message}`); + return null; + }); + if (cached != null) { + return cached; + } + + const summary = await this.loadSummary(filter); + + await this.redisMessagingService + ?.setJsonCache(redisCacheKey, summary, this.redisSummaryCacheMs) + .catch(error => { + console.error(`Share accounting summary cache write failed: ${error.message}`); + }); + + return summary; + } + private async loadSummary(filter: AccountingFilter): Promise { const { whereSql, params } = this.buildWhereClause(filter); const [summary] = await timeAsync('share accounting summary query', () => this.acceptedShareRepository.query(`