Fix pool accounting summary fallback

This commit is contained in:
Ben
2026-07-13 22:47:53 -04:00
parent d830227178
commit 448f7d4c52
4 changed files with 133 additions and 12 deletions
@@ -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);
});
@@ -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<ShareAccountingSummary> {
@@ -380,10 +376,6 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
}
private async withPoolRollupOverlay(summary: ShareAccountingSummary, payoutMode?: PayoutMode): Promise<ShareAccountingSummary> {
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"
+9 -1
View File
@@ -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 { }
@@ -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<void> {
if (this.startupTimer != null) {
clearTimeout(this.startupTimer);
this.startupTimer = null;
}
if (this.timer != null) {
clearInterval(this.timer);
this.timer = null;
}
}
private async refresh(): Promise<void> {
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;
}
}