mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 09:05:06 -07:00
Cache pool accounting summary
This commit is contained in:
@@ -3,6 +3,7 @@ import { InjectRepository } from '@nestjs/typeorm';
|
|||||||
import { Repository } from 'typeorm';
|
import { Repository } from 'typeorm';
|
||||||
|
|
||||||
import { AcceptedShareEntity } from '../accepted-share/accepted-share.entity';
|
import { AcceptedShareEntity } from '../accepted-share/accepted-share.entity';
|
||||||
|
import { RedisMessagingService } from '../../services/redis-messaging.service';
|
||||||
|
|
||||||
export interface AcceptedShareRecord {
|
export interface AcceptedShareRecord {
|
||||||
protocol: 'sv1' | 'sv2';
|
protocol: 'sv1' | 'sv2';
|
||||||
@@ -72,6 +73,7 @@ export class ShareAccountingService implements OnModuleDestroy {
|
|||||||
private flushTimer: NodeJS.Timeout | null = null;
|
private flushTimer: NodeJS.Timeout | null = null;
|
||||||
private activeFlush: Promise<void> | null = null;
|
private activeFlush: Promise<void> | null = null;
|
||||||
private summaryCache = new Map<string, SummaryCacheEntry>();
|
private summaryCache = new Map<string, SummaryCacheEntry>();
|
||||||
|
private readonly poolSummaryCacheKey = 'accounting:pool-summary';
|
||||||
private readonly batchSize = this.readPositiveInt('SHARE_ACCOUNTING_BATCH_SIZE', DEFAULT_BATCH_SIZE);
|
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);
|
private readonly flushIntervalMs = this.readPositiveInt('SHARE_ACCOUNTING_FLUSH_INTERVAL_MS', DEFAULT_FLUSH_INTERVAL_MS);
|
||||||
private readonly maxQueueSize = this.readPositiveInt('SHARE_ACCOUNTING_MAX_QUEUE_SIZE', DEFAULT_MAX_QUEUE_SIZE);
|
private readonly maxQueueSize = this.readPositiveInt('SHARE_ACCOUNTING_MAX_QUEUE_SIZE', DEFAULT_MAX_QUEUE_SIZE);
|
||||||
@@ -81,6 +83,7 @@ export class ShareAccountingService implements OnModuleDestroy {
|
|||||||
constructor(
|
constructor(
|
||||||
@InjectRepository(AcceptedShareEntity)
|
@InjectRepository(AcceptedShareEntity)
|
||||||
private readonly acceptedShareRepository: Repository<AcceptedShareEntity>,
|
private readonly acceptedShareRepository: Repository<AcceptedShareEntity>,
|
||||||
|
private readonly redisMessagingService?: RedisMessagingService,
|
||||||
) { }
|
) { }
|
||||||
|
|
||||||
public async recordAcceptedShare(record: AcceptedShareRecord): Promise<AcceptedShareEntity> {
|
public async recordAcceptedShare(record: AcceptedShareRecord): Promise<AcceptedShareEntity> {
|
||||||
@@ -127,9 +130,29 @@ export class ShareAccountingService implements OnModuleDestroy {
|
|||||||
}
|
}
|
||||||
|
|
||||||
public async getPoolSummary(): Promise<ShareAccountingSummary> {
|
public async getPoolSummary(): Promise<ShareAccountingSummary> {
|
||||||
|
const cached = await this.redisMessagingService
|
||||||
|
?.getJsonCache<ShareAccountingSummary>(this.poolSummaryCacheKey)
|
||||||
|
.catch(error => {
|
||||||
|
console.error(`Pool accounting summary cache read failed: ${error.message}`);
|
||||||
|
return null;
|
||||||
|
});
|
||||||
|
if (cached != null) {
|
||||||
|
return cached;
|
||||||
|
}
|
||||||
|
|
||||||
return this.getSummary({});
|
return this.getSummary({});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public async refreshPoolSummary(): Promise<ShareAccountingSummary> {
|
||||||
|
const summary = await this.getSummary({});
|
||||||
|
await this.redisMessagingService
|
||||||
|
?.setJsonCache(this.poolSummaryCacheKey, summary, 60 * 1000)
|
||||||
|
.catch(error => {
|
||||||
|
console.error(`Pool accounting summary cache write failed: ${error.message}`);
|
||||||
|
});
|
||||||
|
return summary;
|
||||||
|
}
|
||||||
|
|
||||||
public async getAddressSummary(address: string): Promise<ShareAccountingSummary> {
|
public async getAddressSummary(address: string): Promise<ShareAccountingSummary> {
|
||||||
return this.getSummary({ address });
|
return this.getSummary({ address });
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -3,15 +3,18 @@ import { Injectable, OnModuleInit } from '@nestjs/common';
|
|||||||
import { UserAgentReportService } from '../ORM/_views/user-agent-report/user-agent-report.service';
|
import { UserAgentReportService } from '../ORM/_views/user-agent-report/user-agent-report.service';
|
||||||
import { ClientService } from '../ORM/client/client.service';
|
import { ClientService } from '../ORM/client/client.service';
|
||||||
import { RpcBlockService } from '../ORM/rpc-block/rpc-block.service';
|
import { RpcBlockService } from '../ORM/rpc-block/rpc-block.service';
|
||||||
|
import { ShareAccountingService } from '../ORM/share-accounting/share-accounting.service';
|
||||||
|
|
||||||
@Injectable()
|
@Injectable()
|
||||||
export class AppService implements OnModuleInit {
|
export class AppService implements OnModuleInit {
|
||||||
private refreshingLiveUserAgentReport = false;
|
private refreshingLiveUserAgentReport = false;
|
||||||
|
private refreshingPoolSummary = false;
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
private readonly clientService: ClientService,
|
private readonly clientService: ClientService,
|
||||||
private readonly rpcBlockService: RpcBlockService,
|
private readonly rpcBlockService: RpcBlockService,
|
||||||
private readonly userAgentReportService: UserAgentReportService
|
private readonly userAgentReportService: UserAgentReportService,
|
||||||
|
private readonly shareAccountingService: ShareAccountingService
|
||||||
) {
|
) {
|
||||||
|
|
||||||
}
|
}
|
||||||
@@ -30,10 +33,12 @@ export class AppService implements OnModuleInit {
|
|||||||
|
|
||||||
setInterval(async () => {
|
setInterval(async () => {
|
||||||
await this.refreshLiveUserAgentReport();
|
await this.refreshLiveUserAgentReport();
|
||||||
|
await this.refreshPoolSummary();
|
||||||
}, 1000 * 30);
|
}, 1000 * 30);
|
||||||
|
|
||||||
setTimeout(async () => {
|
setTimeout(async () => {
|
||||||
await this.refreshLiveUserAgentReport();
|
await this.refreshLiveUserAgentReport();
|
||||||
|
await this.refreshPoolSummary();
|
||||||
}, 1000 * 15);
|
}, 1000 * 15);
|
||||||
|
|
||||||
setInterval(async () => {
|
setInterval(async () => {
|
||||||
@@ -67,4 +72,19 @@ export class AppService implements OnModuleInit {
|
|||||||
this.refreshingLiveUserAgentReport = false;
|
this.refreshingLiveUserAgentReport = false;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private async refreshPoolSummary() {
|
||||||
|
if (this.refreshingPoolSummary) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
this.refreshingPoolSummary = true;
|
||||||
|
try {
|
||||||
|
await this.shareAccountingService.refreshPoolSummary();
|
||||||
|
} catch (error) {
|
||||||
|
console.error(`Failed refreshing pool accounting summary: ${error.message}`);
|
||||||
|
} finally {
|
||||||
|
this.refreshingPoolSummary = false;
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user