From 4150d1e4256b446db53670a2f6764bb310a41f53 Mon Sep 17 00:00:00 2001 From: Ben Date: Wed, 29 Jul 2026 00:24:48 -0400 Subject: [PATCH] Add pool-level share summary aggregate --- ...SummaryContinuousAggregate1781400000000.ts | 50 +++++++++++++ .../share-accounting.service.spec.ts | 7 +- .../share-accounting.service.ts | 71 ++++++++++++++++++- src/database.config.ts | 2 + test/timescale-redis.integration-spec.ts | 14 +++- 5 files changed, 141 insertions(+), 3 deletions(-) create mode 100644 src/ORM/_migrations/PoolSummaryContinuousAggregate1781400000000.ts diff --git a/src/ORM/_migrations/PoolSummaryContinuousAggregate1781400000000.ts b/src/ORM/_migrations/PoolSummaryContinuousAggregate1781400000000.ts new file mode 100644 index 0000000..12d168e --- /dev/null +++ b/src/ORM/_migrations/PoolSummaryContinuousAggregate1781400000000.ts @@ -0,0 +1,50 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +export class PoolSummaryContinuousAggregate1781400000000 implements MigrationInterface { + public name = 'PoolSummaryContinuousAggregate1781400000000'; + public transaction = false; + + public async up(queryRunner: QueryRunner): Promise { + await queryRunner.query(` + CREATE MATERIALIZED VIEW IF NOT EXISTS "accepted_share_pool_10m" + WITH (timescaledb.continuous) AS + SELECT + time_bucket(INTERVAL '10 minutes', "acceptedAt") AS "bucket", + "payoutMode", + SUM("creditedDifficulty") AS "shares", + COUNT(*) AS "acceptedCount" + FROM "accepted_share_entity" + GROUP BY "bucket", "payoutMode" + WITH NO DATA + `); + await queryRunner.query(` + SELECT add_continuous_aggregate_policy( + 'accepted_share_pool_10m', + start_offset => INTERVAL '2 days', + end_offset => INTERVAL '1 minute', + schedule_interval => INTERVAL '1 minute', + if_not_exists => TRUE + ) + `); + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS "IDX_accepted_share_pool_10m_bucket" + ON "accepted_share_pool_10m" ("bucket" DESC) + `); + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS "IDX_accepted_share_pool_10m_mode_bucket" + ON "accepted_share_pool_10m" ("payoutMode", "bucket" DESC) + `); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_pool_10m_mode_bucket"`); + await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_pool_10m_bucket"`); + await queryRunner.query(` + SELECT remove_continuous_aggregate_policy( + 'accepted_share_pool_10m', + if_exists => TRUE + ) + `); + await queryRunner.query(`DROP MATERIALIZED VIEW IF EXISTS "accepted_share_pool_10m"`); + } +} diff --git a/src/ORM/share-accounting/share-accounting.service.spec.ts b/src/ORM/share-accounting/share-accounting.service.spec.ts index 8c47421..6d6bbaa 100644 --- a/src/ORM/share-accounting/share-accounting.service.spec.ts +++ b/src/ORM/share-accounting/share-accounting.service.spec.ts @@ -318,7 +318,7 @@ describe('ShareAccountingService', () => { expect(repository.query).toHaveBeenNthCalledWith( 1, - expect.stringContaining('FROM "accepted_share_10m"'), + expect.stringContaining('FROM "accepted_share_pool_10m"'), ['solo'], ); expect(repository.query).toHaveBeenNthCalledWith( @@ -413,6 +413,11 @@ describe('ShareAccountingService', () => { networkDifficultyPercent: 35.2, latestShareAt: '2026-06-07T12:10:00.000Z', })); + expect(repository.query).toHaveBeenNthCalledWith( + 1, + expect.stringContaining('FROM "accepted_share_pool_10m"'), + [], + ); expect(repository.query).toHaveBeenNthCalledWith( 2, expect.stringContaining('FROM "accepted_share_block_10m"'), diff --git a/src/ORM/share-accounting/share-accounting.service.ts b/src/ORM/share-accounting/share-accounting.service.ts index 59fc0ec..0cb6723 100644 --- a/src/ORM/share-accounting/share-accounting.service.ts +++ b/src/ORM/share-accounting/share-accounting.service.ts @@ -605,6 +605,16 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy { } private async loadSummary(filter: AccountingFilter): Promise { + if (this.isPoolSummaryFilter(filter)) { + try { + return await this.loadPoolSummary(filter.payoutMode); + } catch (error) { + if (error?.code !== '42P01') { + throw error; + } + } + } + const { whereSql, params } = this.buildWhereClause(filter); const [summary] = await this.acceptedShareRepository.query(` WITH clock AS MATERIALIZED ( @@ -649,6 +659,59 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy { FROM filtered_rows `, params); + return this.mapSummaryRow(summary); + } + + private async loadPoolSummary(payoutMode?: PayoutMode): Promise { + const params = payoutMode == null ? [] : [payoutMode]; + const whereSql = payoutMode == null + ? 'WHERE "accepted_share_pool_10m"."bucket" <= bounds."latestCompletedBucket"' + : 'WHERE "accepted_share_pool_10m"."payoutMode" = $1 AND "accepted_share_pool_10m"."bucket" <= bounds."latestCompletedBucket"'; + const [summary] = await this.acceptedShareRepository.query(` + WITH clock AS MATERIALIZED ( + SELECT + time_bucket(INTERVAL '10 minutes', NOW()) AS "currentBucket" + ), + latest_bucket AS MATERIALIZED ( + SELECT "bucket" + FROM "accepted_share_pool_10m", clock + WHERE "bucket" < clock."currentBucket" + ORDER BY "bucket" DESC + LIMIT 1 + ), + bounds AS MATERIALIZED ( + SELECT + "currentBucket", + COALESCE( + (SELECT "bucket" FROM latest_bucket), + "currentBucket" - INTERVAL '10 minutes' + ) AS "latestCompletedBucket" + FROM clock + ), + filtered_rows AS ( + SELECT "accepted_share_pool_10m".*, bounds."currentBucket", bounds."latestCompletedBucket" + FROM "accepted_share_pool_10m", bounds + ${whereSql} + ) + SELECT + COALESCE(SUM("acceptedCount"), 0)::int AS "totalAcceptedShares", + COALESCE(SUM("shares"), 0)::float AS "totalCreditedDifficulty", + COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" = "latestCompletedBucket"), 0)::int AS "acceptedSharesLast10Minutes", + COALESCE(SUM("shares") FILTER (WHERE "bucket" = "latestCompletedBucket"), 0)::float AS "creditedDifficultyLast10Minutes", + COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 hour' AND "bucket" <= "latestCompletedBucket"), 0)::int AS "acceptedSharesLastHour", + COALESCE(SUM("shares") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 hour' AND "bucket" <= "latestCompletedBucket"), 0)::float AS "creditedDifficultyLastHour", + COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 day' AND "bucket" <= "latestCompletedBucket"), 0)::int AS "acceptedSharesLastDay", + COALESCE(SUM("shares") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 day' AND "bucket" <= "latestCompletedBucket"), 0)::float AS "creditedDifficultyLastDay", + COALESCE((SUM("shares") FILTER (WHERE "bucket" = "latestCompletedBucket") * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes", + COALESCE((SUM("shares") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 hour' AND "bucket" <= "latestCompletedBucket") * ${HASHES_PER_DIFFICULTY}) / 3600, 0)::float AS "hashRateLastHour", + MAX("bucket") AS "latestShareAt" + FROM filtered_rows + `, params); + + return this.mapSummaryRow(summary); + } + + private mapSummaryRow(summary: Record | undefined): ShareAccountingSummary { return { totalAcceptedShares: this.toNumber(summary?.totalAcceptedShares), totalCreditedDifficulty: this.toNumber(summary?.totalCreditedDifficulty), @@ -669,11 +732,17 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy { blockCandidateCount: 0, latestShareAt: summary?.latestShareAt == null ? null - : new Date(summary.latestShareAt).toISOString(), + : new Date(summary.latestShareAt as string | Date).toISOString(), protocolBreakdown: [], }; } + private isPoolSummaryFilter(filter: AccountingFilter): boolean { + return filter.address == null + && filter.clientName == null + && filter.clientId == null; + } + private scheduleFlush(): void { if (this.flushTimer != null || this.activeFlush != null) { return; diff --git a/src/database.config.ts b/src/database.config.ts index 8a1de75..a10d956 100644 --- a/src/database.config.ts +++ b/src/database.config.ts @@ -18,6 +18,7 @@ import { PayoutModes1781300000000 } from './ORM/_migrations/PayoutModes178130000 import { ShareRollupStoragePolicy1781305000000 } from './ORM/_migrations/ShareRollupStoragePolicy1781305000000'; import { AcceptedShareHighScores1781309000000 } from './ORM/_migrations/AcceptedShareHighScores1781309000000'; import { UserAgentReportNonzeroHashrate1781313000000 } from './ORM/_migrations/UserAgentReportNonzeroHashrate1781313000000'; +import { PoolSummaryContinuousAggregate1781400000000 } from './ORM/_migrations/PoolSummaryContinuousAggregate1781400000000'; 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'; @@ -62,6 +63,7 @@ export const databaseMigrations = [ ShareRollupStoragePolicy1781305000000, AcceptedShareHighScores1781309000000, UserAgentReportNonzeroHashrate1781313000000, + PoolSummaryContinuousAggregate1781400000000, ]; export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions { diff --git a/test/timescale-redis.integration-spec.ts b/test/timescale-redis.integration-spec.ts index d0004ec..9fc9c86 100644 --- a/test/timescale-redis.integration-spec.ts +++ b/test/timescale-redis.integration-spec.ts @@ -76,7 +76,7 @@ describe('TimescaleDB and Redis integration', () => { const aggregates = await dataSource.query(` SELECT view_name FROM timescaledb_information.continuous_aggregates - WHERE view_name IN ('accepted_share_10m', 'accepted_share_1h', 'accepted_share_1d', 'accepted_share_block_10m') + WHERE view_name IN ('accepted_share_10m', 'accepted_share_1h', 'accepted_share_1d', 'accepted_share_block_10m', 'accepted_share_pool_10m') ORDER BY view_name `); expect(aggregates.map(row => row.view_name)).toEqual([ @@ -84,6 +84,7 @@ describe('TimescaleDB and Redis integration', () => { 'accepted_share_1d', 'accepted_share_1h', 'accepted_share_block_10m', + 'accepted_share_pool_10m', ]); const legacyTables = await dataSource.query(` @@ -228,6 +229,7 @@ describe('TimescaleDB and Redis integration', () => { await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_10m', NULL, NULL)`); await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_block_10m', NULL, NULL)`); + await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_pool_10m', NULL, NULL)`); const aggregateRows = await dataSource.query(` SELECT "shares"::float AS shares, "acceptedCount"::int AS "acceptedCount" FROM accepted_share_10m @@ -250,6 +252,16 @@ describe('TimescaleDB and Redis integration', () => { expect(blockAggregateRows).toEqual(expect.arrayContaining([ expect.objectContaining({ shares: 96, acceptedCount: 2, networkDifficulty: 100000 }), ])); + + const poolAggregateRows = await dataSource.query(` + SELECT "shares"::float AS shares, "acceptedCount"::int AS "acceptedCount" + FROM accepted_share_pool_10m + WHERE "payoutMode" = $1 + `, ['solo']); + + expect(poolAggregateRows).toEqual(expect.arrayContaining([ + expect.objectContaining({ shares: 96, acceptedCount: 2 }), + ])); }); it('should retain daily and all-time best share from completed rollup buckets', async () => {