From 448f7d4c527ab774d1e901a9f01f714c5c34e5af Mon Sep 17 00:00:00 2001 From: Ben Date: Mon, 13 Jul 2026 22:47:53 -0400 Subject: [PATCH] Fix pool accounting summary fallback --- .../share-accounting.service.spec.ts | 60 ++++++++++++++++- .../share-accounting.service.ts | 10 +-- src/notifier.module.ts | 10 ++- src/services/pool-summary-refresh.service.ts | 65 +++++++++++++++++++ 4 files changed, 133 insertions(+), 12 deletions(-) create mode 100644 src/services/pool-summary-refresh.service.ts diff --git a/src/ORM/share-accounting/share-accounting.service.spec.ts b/src/ORM/share-accounting/share-accounting.service.spec.ts index f8a6066..8c47421 100644 --- a/src/ORM/share-accounting/share-accounting.service.spec.ts +++ b/src/ORM/share-accounting/share-accounting.service.spec.ts @@ -272,6 +272,62 @@ describe('ShareAccountingService', () => { ); }); + it('should serve API-only pool accounting from rollups when the precomputed pool cache is missing', async () => { + process.env.API_ONLY = 'true'; + const redis = { + getJsonCache: jest.fn() + .mockResolvedValueOnce(null) + .mockResolvedValueOnce(null), + setJsonCache: jest.fn().mockResolvedValue(undefined), + }; + 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'), + }]) + .mockResolvedValueOnce([{ + bestSubmissionDifficulty: '4096', + bestSubmissionDifficultyAt: new Date('2026-06-07T12:20:00Z'), + currentRoundAcceptedShares: '11', + workSinceLastBlock: '352', + currentRoundNetworkDifficulty: '1000', + }]) + .mockResolvedValueOnce([]), + }; + const service = new ShareAccountingService(repository as any, redis as any); + + await expect(service.getPoolSummary('solo')).resolves.toEqual(expect.objectContaining({ + totalAcceptedShares: 12, + totalCreditedDifficulty: 384, + acceptedSharesLast10Minutes: 4, + creditedDifficultyLast10Minutes: 128, + bestSubmissionDifficulty: 4096, + workSinceLastBlock: 352, + networkDifficultyPercent: 35.2, + })); + + expect(repository.query).toHaveBeenNthCalledWith( + 1, + expect.stringContaining('FROM "accepted_share_10m"'), + ['solo'], + ); + expect(repository.query).toHaveBeenNthCalledWith( + 2, + expect.stringContaining('FROM "accepted_share_block_10m"'), + ['solo'], + ); + }); + it('should use Redis cache for share accounting summaries across API workers', async () => { process.env.SHARE_ACCOUNTING_SUMMARY_CACHE_MS = '0'; const repository = { @@ -429,8 +485,8 @@ describe('ShareAccountingService', () => { }; const service = new ShareAccountingService(repository as any); - await service.getPoolSummary(); - await service.getPoolSummary(); + await service.getAddressSummary('bc1qcached'); + await service.getAddressSummary('bc1qcached'); expect(repository.query).toHaveBeenCalledTimes(1); }); diff --git a/src/ORM/share-accounting/share-accounting.service.ts b/src/ORM/share-accounting/share-accounting.service.ts index fe55cde..59fc0ec 100644 --- a/src/ORM/share-accounting/share-accounting.service.ts +++ b/src/ORM/share-accounting/share-accounting.service.ts @@ -360,11 +360,7 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy { return cached; } - if (process.env.API_ONLY === 'true') { - return this.emptySummary(); - } - - return this.getSummary({ payoutMode: mode }); + return this.withPoolRollupOverlay(await this.getSummary({ payoutMode: mode }), mode); } public async refreshPoolSummary(payoutMode?: PayoutMode): Promise { @@ -380,10 +376,6 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy { } private async withPoolRollupOverlay(summary: ShareAccountingSummary, payoutMode?: PayoutMode): Promise { - if (process.env.API_ONLY === 'true') { - return summary; - } - const [currentRoundRow] = await this.acceptedShareRepository.query(` WITH latest_found_block AS ( SELECT COALESCE(MAX("height"), 0) AS "height" diff --git a/src/notifier.module.ts b/src/notifier.module.ts index bacf5cd..ba8b29f 100644 --- a/src/notifier.module.ts +++ b/src/notifier.module.ts @@ -3,9 +3,12 @@ import { ConfigModule, ConfigService } from '@nestjs/config'; import { TypeOrmModule } from '@nestjs/typeorm'; import { createDatabaseOptions } from './database.config'; +import { AcceptedShareEntity } from './ORM/accepted-share/accepted-share.entity'; import { PayoutSnapshotModule } from './ORM/payout-snapshot/payout-snapshot.module'; import { RpcBlocksModule } from './ORM/rpc-block/rpc-block.module'; +import { ShareAccountingService } from './ORM/share-accounting/share-accounting.service'; import { BitcoinRpcService } from './services/bitcoin-rpc.service'; +import { PoolSummaryRefreshService } from './services/pool-summary-refresh.service'; import { RedisMessagingModule } from './services/redis-messaging.module'; /** @@ -30,10 +33,15 @@ import { RedisMessagingModule } from './services/redis-messaging.module'; DB_POOL_SIZE: configService.get('DB_POOL_SIZE'), }), }), + TypeOrmModule.forFeature([AcceptedShareEntity]), RedisMessagingModule, RpcBlocksModule, PayoutSnapshotModule, ], - providers: [BitcoinRpcService], + providers: [ + BitcoinRpcService, + PoolSummaryRefreshService, + ShareAccountingService, + ], }) export class NotifierModule { } diff --git a/src/services/pool-summary-refresh.service.ts b/src/services/pool-summary-refresh.service.ts new file mode 100644 index 0000000..bfbf1df --- /dev/null +++ b/src/services/pool-summary-refresh.service.ts @@ -0,0 +1,65 @@ +import { Injectable, OnModuleDestroy, OnModuleInit } from '@nestjs/common'; + +import { ShareAccountingService } from '../ORM/share-accounting/share-accounting.service'; + +const DEFAULT_REFRESH_INTERVAL_MS = 5 * 60 * 1000; +const DEFAULT_STARTUP_DELAY_MS = 15 * 1000; + +@Injectable() +export class PoolSummaryRefreshService implements OnModuleInit, OnModuleDestroy { + private timer: NodeJS.Timeout | null = null; + private startupTimer: NodeJS.Timeout | null = null; + private refreshing = false; + + constructor(private readonly shareAccountingService: ShareAccountingService) { } + + public onModuleInit(): void { + if (process.env.MASTER !== 'true' || process.env.API_ONLY === 'true') { + return; + } + + this.startupTimer = setTimeout(() => { + void this.refresh(); + }, this.readPositiveInt('POOL_SUMMARY_REFRESH_STARTUP_DELAY_MS', DEFAULT_STARTUP_DELAY_MS)); + this.startupTimer.unref?.(); + + this.timer = setInterval(() => { + void this.refresh(); + }, this.readPositiveInt('POOL_SUMMARY_REFRESH_INTERVAL_MS', DEFAULT_REFRESH_INTERVAL_MS)); + this.timer.unref?.(); + } + + public async onModuleDestroy(): Promise { + if (this.startupTimer != null) { + clearTimeout(this.startupTimer); + this.startupTimer = null; + } + + if (this.timer != null) { + clearInterval(this.timer); + this.timer = null; + } + } + + private async refresh(): Promise { + if (this.refreshing) { + return; + } + + this.refreshing = true; + try { + await this.shareAccountingService.refreshPoolSummary(); + await this.shareAccountingService.refreshPoolSummary('pplns'); + await this.shareAccountingService.refreshPoolSummary('solo'); + } catch (error) { + console.error(`Failed refreshing pool accounting summary: ${error.message}`); + } finally { + this.refreshing = false; + } + } + + private readPositiveInt(name: string, fallback: number): number { + const value = Number(process.env[name]); + return Number.isInteger(value) && value > 0 ? value : fallback; + } +}