Serve client accounting from rollup cache

This commit is contained in:
Ben
2026-06-08 01:42:17 -04:00
parent 9b73e33b13
commit deaf3e56bd
2 changed files with 107 additions and 5 deletions
@@ -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 = {
@@ -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<ShareAccountingSummary> {
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<ShareAccountingSummary> {
const redisCacheKey = `accounting:summary:${cacheKey}`;
const cached = await this.redisMessagingService
?.getJsonCache<ShareAccountingSummary>(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<ShareAccountingSummary> {
const { whereSql, params } = this.buildWhereClause(filter);
const [summary] = await timeAsync('share accounting summary query', () => this.acceptedShareRepository.query(`