From 4c91e69ca88c5349383a6e7128e82810c11acdae Mon Sep 17 00:00:00 2001 From: Ben Date: Mon, 8 Jun 2026 00:55:51 -0400 Subject: [PATCH] Use share rollups for API accounting --- ...AcceptedShareRollupIndexes1780865400000.ts | 32 +++++++++++ .../share-accounting.service.spec.ts | 34 +++-------- .../share-accounting.service.ts | 57 +++++++------------ src/database.config.ts | 2 + 4 files changed, 62 insertions(+), 63 deletions(-) create mode 100644 src/ORM/_migrations/AcceptedShareRollupIndexes1780865400000.ts diff --git a/src/ORM/_migrations/AcceptedShareRollupIndexes1780865400000.ts b/src/ORM/_migrations/AcceptedShareRollupIndexes1780865400000.ts new file mode 100644 index 0000000..463b6fc --- /dev/null +++ b/src/ORM/_migrations/AcceptedShareRollupIndexes1780865400000.ts @@ -0,0 +1,32 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +export class AcceptedShareRollupIndexes1780865400000 implements MigrationInterface { + public name = 'AcceptedShareRollupIndexes1780865400000'; + public transaction = false; + + public async up(queryRunner: QueryRunner): Promise { + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_bucket" + ON "accepted_share_10m" ("bucket" DESC) + `); + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_address_bucket" + ON "accepted_share_10m" ("address", "bucket" DESC) + `); + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_group_bucket" + ON "accepted_share_10m" ("address", "clientName", "bucket" DESC) + `); + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_client_bucket" + ON "accepted_share_10m" ("clientId", "bucket" DESC) + `); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_client_bucket"`); + await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_group_bucket"`); + await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_address_bucket"`); + await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_bucket"`); + } +} diff --git a/src/ORM/share-accounting/share-accounting.service.spec.ts b/src/ORM/share-accounting/share-accounting.service.spec.ts index f101491..67331a8 100644 --- a/src/ORM/share-accounting/share-accounting.service.spec.ts +++ b/src/ORM/share-accounting/share-accounting.service.spec.ts @@ -121,7 +121,7 @@ describe('ShareAccountingService', () => { await expect(queued).resolves.toEqual(expect.objectContaining({ jobId: 'queued' })); }); - it('should return numeric accounting summaries with protocol breakdown', async () => { + it('should return numeric accounting summaries from the share rollup', async () => { const repository = { query: jest.fn() .mockResolvedValueOnce([{ @@ -135,14 +135,7 @@ describe('ShareAccountingService', () => { creditedDifficultyLastDay: '96', hashRateLast10Minutes: '458129844.9', hashRateLastHour: '114532461.2', - bestSubmissionDifficulty: '2048', - blockCandidateCount: '1', latestShareAt: new Date('2026-06-07T12:10:00Z'), - }]) - .mockResolvedValueOnce([{ - protocol: 'sv1', - acceptedShares: '3', - creditedDifficulty: '96', }]), }; const service = new ShareAccountingService(repository as any); @@ -158,25 +151,17 @@ describe('ShareAccountingService', () => { creditedDifficultyLastDay: 96, hashRateLast10Minutes: 458129844.9, hashRateLastHour: 114532461.2, - bestSubmissionDifficulty: 2048, - blockCandidateCount: 1, + bestSubmissionDifficulty: 0, + blockCandidateCount: 0, latestShareAt: '2026-06-07T12:10:00.000Z', - protocolBreakdown: [{ - protocol: 'sv1', - acceptedShares: 3, - creditedDifficulty: 96, - }], + protocolBreakdown: [], }); expect(repository.query).toHaveBeenNthCalledWith( 1, - expect.stringContaining('"address" = $1'), - ['bc1qtest'], - ); - expect(repository.query).toHaveBeenNthCalledWith( - 2, - expect.stringContaining('GROUP BY "protocol"'), + expect.stringContaining('"accepted_share_10m"'), ['bc1qtest'], ); + expect(repository.query).toHaveBeenCalledTimes(1); }); it('should cache accounting summaries briefly to protect hot dashboard endpoints', async () => { @@ -194,18 +179,15 @@ describe('ShareAccountingService', () => { creditedDifficultyLastDay: '32', hashRateLast10Minutes: '1', hashRateLastHour: '1', - bestSubmissionDifficulty: '32', - blockCandidateCount: '0', latestShareAt: null, - }]) - .mockResolvedValueOnce([{ protocol: 'sv1', acceptedShares: '1', creditedDifficulty: '32' }]), + }]), }; const service = new ShareAccountingService(repository as any); await service.getPoolSummary(); await service.getPoolSummary(); - expect(repository.query).toHaveBeenCalledTimes(2); + 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 a9558f6..542e8c6 100644 --- a/src/ORM/share-accounting/share-accounting.service.ts +++ b/src/ORM/share-accounting/share-accounting.service.ts @@ -62,6 +62,7 @@ interface AccountingFilter { } const HASHES_PER_DIFFICULTY = 4294967296; +const ROLLUP_BUCKET_SECONDS = 600; const DEFAULT_BATCH_SIZE = 500; const DEFAULT_FLUSH_INTERVAL_MS = 25; const DEFAULT_MAX_QUEUE_SIZE = 50000; @@ -199,10 +200,9 @@ export class ShareAccountingService implements OnModuleDestroy { const rows = await timeAsync('share accounting session summaries query', () => this.acceptedShareRepository.query(` SELECT "clientId", - MAX("acceptedAt") AS "latestShareAt", - COALESCE((SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes') * ${HASHES_PER_DIFFICULTY}) / 600, 0)::float AS "hashRateLast10Minutes", - COALESCE(MAX("submissionDifficulty"), 0)::float AS "bestSubmissionDifficulty" - FROM "accepted_share_entity" + MAX("bucket") AS "latestShareAt", + COALESCE((SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes') * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes" + FROM "accepted_share_10m" WHERE "clientId" = ANY($1::uuid[]) GROUP BY "clientId" `, [uniqueClientIds]), { clientIds: uniqueClientIds.length }); @@ -214,7 +214,7 @@ export class ShareAccountingService implements OnModuleDestroy { ? null : new Date(row.latestShareAt).toISOString(), hashRateLast10Minutes: this.toNumber(row.hashRateLast10Minutes), - bestSubmissionDifficulty: this.toNumber(row.bestSubmissionDifficulty), + bestSubmissionDifficulty: 0, }); }); @@ -250,34 +250,21 @@ export class ShareAccountingService implements OnModuleDestroy { const { whereSql, params } = this.buildWhereClause(filter); const [summary] = await timeAsync('share accounting summary query', () => this.acceptedShareRepository.query(` SELECT - COUNT(*)::int AS "totalAcceptedShares", - COALESCE(SUM("creditedDifficulty"), 0)::float AS "totalCreditedDifficulty", - COUNT(*) FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes')::int AS "acceptedSharesLast10Minutes", - COALESCE(SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes'), 0)::float AS "creditedDifficultyLast10Minutes", - COUNT(*) FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 hour')::int AS "acceptedSharesLastHour", - COALESCE(SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 hour'), 0)::float AS "creditedDifficultyLastHour", - COUNT(*) FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 day')::int AS "acceptedSharesLastDay", - COALESCE(SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 day'), 0)::float AS "creditedDifficultyLastDay", - COALESCE((SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes') * ${HASHES_PER_DIFFICULTY}) / 600, 0)::float AS "hashRateLast10Minutes", - COALESCE((SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 hour') * ${HASHES_PER_DIFFICULTY}) / 3600, 0)::float AS "hashRateLastHour", - COALESCE(MAX("submissionDifficulty"), 0)::float AS "bestSubmissionDifficulty", - COUNT(*) FILTER (WHERE "isBlockCandidate" = TRUE)::int AS "blockCandidateCount", - MAX("acceptedAt") AS "latestShareAt" - FROM "accepted_share_entity" + COALESCE(SUM("acceptedCount"), 0)::int AS "totalAcceptedShares", + COALESCE(SUM("shares"), 0)::float AS "totalCreditedDifficulty", + COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes'), 0)::int AS "acceptedSharesLast10Minutes", + COALESCE(SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes'), 0)::float AS "creditedDifficultyLast10Minutes", + COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 hour'), 0)::int AS "acceptedSharesLastHour", + COALESCE(SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 hour'), 0)::float AS "creditedDifficultyLastHour", + COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 day'), 0)::int AS "acceptedSharesLastDay", + COALESCE(SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 day'), 0)::float AS "creditedDifficultyLastDay", + COALESCE((SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes') * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes", + COALESCE((SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 hour') * ${HASHES_PER_DIFFICULTY}) / 3600, 0)::float AS "hashRateLastHour", + MAX("bucket") AS "latestShareAt" + FROM "accepted_share_10m" ${whereSql} `, params), { filter, params: params.length }); - const protocolRows = await timeAsync('share accounting protocol breakdown query', () => this.acceptedShareRepository.query(` - SELECT - "protocol", - COUNT(*)::int AS "acceptedShares", - COALESCE(SUM("creditedDifficulty"), 0)::float AS "creditedDifficulty" - FROM "accepted_share_entity" - ${whereSql} - GROUP BY "protocol" - ORDER BY "protocol" - `, params), { filter, params: params.length }); - return { totalAcceptedShares: this.toNumber(summary?.totalAcceptedShares), totalCreditedDifficulty: this.toNumber(summary?.totalCreditedDifficulty), @@ -289,16 +276,12 @@ export class ShareAccountingService implements OnModuleDestroy { creditedDifficultyLastDay: this.toNumber(summary?.creditedDifficultyLastDay), hashRateLast10Minutes: this.toNumber(summary?.hashRateLast10Minutes), hashRateLastHour: this.toNumber(summary?.hashRateLastHour), - bestSubmissionDifficulty: this.toNumber(summary?.bestSubmissionDifficulty), - blockCandidateCount: this.toNumber(summary?.blockCandidateCount), + bestSubmissionDifficulty: 0, + blockCandidateCount: 0, latestShareAt: summary?.latestShareAt == null ? null : new Date(summary.latestShareAt).toISOString(), - protocolBreakdown: protocolRows.map(row => ({ - protocol: row.protocol, - acceptedShares: this.toNumber(row.acceptedShares), - creditedDifficulty: this.toNumber(row.creditedDifficulty), - })), + protocolBreakdown: [], }; } diff --git a/src/database.config.ts b/src/database.config.ts index f9033b1..e0d85c2 100644 --- a/src/database.config.ts +++ b/src/database.config.ts @@ -5,6 +5,7 @@ import { InitialTimescaleSchema1780859300000 } from './ORM/_migrations/InitialTi import { ActiveOnlyUserAgentReport1780860200000 } from './ORM/_migrations/ActiveOnlyUserAgentReport1780860200000'; import { TimescaleOperationalHardening1780861200000 } from './ORM/_migrations/TimescaleOperationalHardening1780861200000'; import { AcceptedShareIndex1780862400000 } from './ORM/_migrations/AcceptedShareIndex1780862400000'; +import { AcceptedShareRollupIndexes1780865400000 } from './ORM/_migrations/AcceptedShareRollupIndexes1780865400000'; import { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-report.view'; import { AcceptedShareEntity } from './ORM/accepted-share/accepted-share.entity'; import { AddressSettingsEntity } from './ORM/address-settings/address-settings.entity'; @@ -28,6 +29,7 @@ export const databaseMigrations = [ ActiveOnlyUserAgentReport1780860200000, TimescaleOperationalHardening1780861200000, AcceptedShareIndex1780862400000, + AcceptedShareRollupIndexes1780865400000, ]; export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions {