Use share rollups for API accounting

This commit is contained in:
Ben
2026-06-08 00:55:51 -04:00
parent 7b7dc87242
commit 4c91e69ca8
4 changed files with 62 additions and 63 deletions
@@ -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<void> {
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<void> {
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"`);
}
}
@@ -121,7 +121,7 @@ describe('ShareAccountingService', () => {
await expect(queued).resolves.toEqual(expect.objectContaining({ jobId: 'queued' })); 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 = { const repository = {
query: jest.fn() query: jest.fn()
.mockResolvedValueOnce([{ .mockResolvedValueOnce([{
@@ -135,14 +135,7 @@ describe('ShareAccountingService', () => {
creditedDifficultyLastDay: '96', creditedDifficultyLastDay: '96',
hashRateLast10Minutes: '458129844.9', hashRateLast10Minutes: '458129844.9',
hashRateLastHour: '114532461.2', hashRateLastHour: '114532461.2',
bestSubmissionDifficulty: '2048',
blockCandidateCount: '1',
latestShareAt: new Date('2026-06-07T12:10:00Z'), latestShareAt: new Date('2026-06-07T12:10:00Z'),
}])
.mockResolvedValueOnce([{
protocol: 'sv1',
acceptedShares: '3',
creditedDifficulty: '96',
}]), }]),
}; };
const service = new ShareAccountingService(repository as any); const service = new ShareAccountingService(repository as any);
@@ -158,25 +151,17 @@ describe('ShareAccountingService', () => {
creditedDifficultyLastDay: 96, creditedDifficultyLastDay: 96,
hashRateLast10Minutes: 458129844.9, hashRateLast10Minutes: 458129844.9,
hashRateLastHour: 114532461.2, hashRateLastHour: 114532461.2,
bestSubmissionDifficulty: 2048, bestSubmissionDifficulty: 0,
blockCandidateCount: 1, blockCandidateCount: 0,
latestShareAt: '2026-06-07T12:10:00.000Z', latestShareAt: '2026-06-07T12:10:00.000Z',
protocolBreakdown: [{ protocolBreakdown: [],
protocol: 'sv1',
acceptedShares: 3,
creditedDifficulty: 96,
}],
}); });
expect(repository.query).toHaveBeenNthCalledWith( expect(repository.query).toHaveBeenNthCalledWith(
1, 1,
expect.stringContaining('"address" = $1'), expect.stringContaining('"accepted_share_10m"'),
['bc1qtest'],
);
expect(repository.query).toHaveBeenNthCalledWith(
2,
expect.stringContaining('GROUP BY "protocol"'),
['bc1qtest'], ['bc1qtest'],
); );
expect(repository.query).toHaveBeenCalledTimes(1);
}); });
it('should cache accounting summaries briefly to protect hot dashboard endpoints', async () => { it('should cache accounting summaries briefly to protect hot dashboard endpoints', async () => {
@@ -194,18 +179,15 @@ describe('ShareAccountingService', () => {
creditedDifficultyLastDay: '32', creditedDifficultyLastDay: '32',
hashRateLast10Minutes: '1', hashRateLast10Minutes: '1',
hashRateLastHour: '1', hashRateLastHour: '1',
bestSubmissionDifficulty: '32',
blockCandidateCount: '0',
latestShareAt: null, latestShareAt: null,
}]) }]),
.mockResolvedValueOnce([{ protocol: 'sv1', acceptedShares: '1', creditedDifficulty: '32' }]),
}; };
const service = new ShareAccountingService(repository as any); const service = new ShareAccountingService(repository as any);
await service.getPoolSummary(); await service.getPoolSummary();
await service.getPoolSummary(); await service.getPoolSummary();
expect(repository.query).toHaveBeenCalledTimes(2); expect(repository.query).toHaveBeenCalledTimes(1);
}); });
}); });
@@ -62,6 +62,7 @@ interface AccountingFilter {
} }
const HASHES_PER_DIFFICULTY = 4294967296; const HASHES_PER_DIFFICULTY = 4294967296;
const ROLLUP_BUCKET_SECONDS = 600;
const DEFAULT_BATCH_SIZE = 500; const DEFAULT_BATCH_SIZE = 500;
const DEFAULT_FLUSH_INTERVAL_MS = 25; const DEFAULT_FLUSH_INTERVAL_MS = 25;
const DEFAULT_MAX_QUEUE_SIZE = 50000; 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(` const rows = await timeAsync('share accounting session summaries query', () => this.acceptedShareRepository.query(`
SELECT SELECT
"clientId", "clientId",
MAX("acceptedAt") AS "latestShareAt", MAX("bucket") AS "latestShareAt",
COALESCE((SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes') * ${HASHES_PER_DIFFICULTY}) / 600, 0)::float AS "hashRateLast10Minutes", COALESCE((SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes') * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes"
COALESCE(MAX("submissionDifficulty"), 0)::float AS "bestSubmissionDifficulty" FROM "accepted_share_10m"
FROM "accepted_share_entity"
WHERE "clientId" = ANY($1::uuid[]) WHERE "clientId" = ANY($1::uuid[])
GROUP BY "clientId" GROUP BY "clientId"
`, [uniqueClientIds]), { clientIds: uniqueClientIds.length }); `, [uniqueClientIds]), { clientIds: uniqueClientIds.length });
@@ -214,7 +214,7 @@ export class ShareAccountingService implements OnModuleDestroy {
? null ? null
: new Date(row.latestShareAt).toISOString(), : new Date(row.latestShareAt).toISOString(),
hashRateLast10Minutes: this.toNumber(row.hashRateLast10Minutes), 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 { whereSql, params } = this.buildWhereClause(filter);
const [summary] = await timeAsync('share accounting summary query', () => this.acceptedShareRepository.query(` const [summary] = await timeAsync('share accounting summary query', () => this.acceptedShareRepository.query(`
SELECT SELECT
COUNT(*)::int AS "totalAcceptedShares", COALESCE(SUM("acceptedCount"), 0)::int AS "totalAcceptedShares",
COALESCE(SUM("creditedDifficulty"), 0)::float AS "totalCreditedDifficulty", COALESCE(SUM("shares"), 0)::float AS "totalCreditedDifficulty",
COUNT(*) FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes')::int AS "acceptedSharesLast10Minutes", COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes'), 0)::int AS "acceptedSharesLast10Minutes",
COALESCE(SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes'), 0)::float AS "creditedDifficultyLast10Minutes", COALESCE(SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes'), 0)::float AS "creditedDifficultyLast10Minutes",
COUNT(*) FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 hour')::int AS "acceptedSharesLastHour", COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 hour'), 0)::int AS "acceptedSharesLastHour",
COALESCE(SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 hour'), 0)::float AS "creditedDifficultyLastHour", COALESCE(SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 hour'), 0)::float AS "creditedDifficultyLastHour",
COUNT(*) FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 day')::int AS "acceptedSharesLastDay", COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 day'), 0)::int AS "acceptedSharesLastDay",
COALESCE(SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 day'), 0)::float AS "creditedDifficultyLastDay", COALESCE(SUM("shares") FILTER (WHERE "bucket" > 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("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes') * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes",
COALESCE((SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 hour') * ${HASHES_PER_DIFFICULTY}) / 3600, 0)::float AS "hashRateLastHour", COALESCE((SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 hour') * ${HASHES_PER_DIFFICULTY}) / 3600, 0)::float AS "hashRateLastHour",
COALESCE(MAX("submissionDifficulty"), 0)::float AS "bestSubmissionDifficulty", MAX("bucket") AS "latestShareAt"
COUNT(*) FILTER (WHERE "isBlockCandidate" = TRUE)::int AS "blockCandidateCount", FROM "accepted_share_10m"
MAX("acceptedAt") AS "latestShareAt"
FROM "accepted_share_entity"
${whereSql} ${whereSql}
`, params), { filter, params: params.length }); `, 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 { return {
totalAcceptedShares: this.toNumber(summary?.totalAcceptedShares), totalAcceptedShares: this.toNumber(summary?.totalAcceptedShares),
totalCreditedDifficulty: this.toNumber(summary?.totalCreditedDifficulty), totalCreditedDifficulty: this.toNumber(summary?.totalCreditedDifficulty),
@@ -289,16 +276,12 @@ export class ShareAccountingService implements OnModuleDestroy {
creditedDifficultyLastDay: this.toNumber(summary?.creditedDifficultyLastDay), creditedDifficultyLastDay: this.toNumber(summary?.creditedDifficultyLastDay),
hashRateLast10Minutes: this.toNumber(summary?.hashRateLast10Minutes), hashRateLast10Minutes: this.toNumber(summary?.hashRateLast10Minutes),
hashRateLastHour: this.toNumber(summary?.hashRateLastHour), hashRateLastHour: this.toNumber(summary?.hashRateLastHour),
bestSubmissionDifficulty: this.toNumber(summary?.bestSubmissionDifficulty), bestSubmissionDifficulty: 0,
blockCandidateCount: this.toNumber(summary?.blockCandidateCount), blockCandidateCount: 0,
latestShareAt: summary?.latestShareAt == null latestShareAt: summary?.latestShareAt == null
? null ? null
: new Date(summary.latestShareAt).toISOString(), : new Date(summary.latestShareAt).toISOString(),
protocolBreakdown: protocolRows.map(row => ({ protocolBreakdown: [],
protocol: row.protocol,
acceptedShares: this.toNumber(row.acceptedShares),
creditedDifficulty: this.toNumber(row.creditedDifficulty),
})),
}; };
} }
+2
View File
@@ -5,6 +5,7 @@ import { InitialTimescaleSchema1780859300000 } from './ORM/_migrations/InitialTi
import { ActiveOnlyUserAgentReport1780860200000 } from './ORM/_migrations/ActiveOnlyUserAgentReport1780860200000'; import { ActiveOnlyUserAgentReport1780860200000 } from './ORM/_migrations/ActiveOnlyUserAgentReport1780860200000';
import { TimescaleOperationalHardening1780861200000 } from './ORM/_migrations/TimescaleOperationalHardening1780861200000'; import { TimescaleOperationalHardening1780861200000 } from './ORM/_migrations/TimescaleOperationalHardening1780861200000';
import { AcceptedShareIndex1780862400000 } from './ORM/_migrations/AcceptedShareIndex1780862400000'; 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 { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-report.view';
import { AcceptedShareEntity } from './ORM/accepted-share/accepted-share.entity'; import { AcceptedShareEntity } from './ORM/accepted-share/accepted-share.entity';
import { AddressSettingsEntity } from './ORM/address-settings/address-settings.entity'; import { AddressSettingsEntity } from './ORM/address-settings/address-settings.entity';
@@ -28,6 +29,7 @@ export const databaseMigrations = [
ActiveOnlyUserAgentReport1780860200000, ActiveOnlyUserAgentReport1780860200000,
TimescaleOperationalHardening1780861200000, TimescaleOperationalHardening1780861200000,
AcceptedShareIndex1780862400000, AcceptedShareIndex1780862400000,
AcceptedShareRollupIndexes1780865400000,
]; ];
export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions { export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions {