From bcd6bb0aa7b4c432f574f1a616c1448cc7844aea Mon Sep 17 00:00:00 2001 From: Ben Date: Tue, 9 Jun 2026 22:55:11 -0400 Subject: [PATCH] Add accepted share rollup batches --- .env.example | 4 + docker-compose.external-db.yml | 4 + full-setup/docker-compose-mainnet.yml | 4 + ...dShareRetentionCompression1780966200000.ts | 54 +++++ .../ShareRollupBatches1780962600000.ts | 62 +++++ .../share-accounting.service.spec.ts | 63 +++++ .../share-accounting.service.ts | 224 +++++++++++++++++- src/database.config.ts | 4 + test/timescale-redis.integration-spec.ts | 107 +++++++++ 9 files changed, 524 insertions(+), 2 deletions(-) create mode 100644 src/ORM/_migrations/AcceptedShareRetentionCompression1780966200000.ts create mode 100644 src/ORM/_migrations/ShareRollupBatches1780962600000.ts diff --git a/.env.example b/.env.example index e908fee..60a8abf 100644 --- a/.env.example +++ b/.env.example @@ -77,6 +77,10 @@ SHARE_ACCOUNTING_FLUSH_INTERVAL_MS=25 SHARE_ACCOUNTING_MAX_QUEUE_SIZE=50000 SHARE_ACCOUNTING_SUMMARY_CACHE_MS=2500 SHARE_ACCOUNTING_SUMMARY_CACHE_MAX=10000 +SHARE_ROLLUP_ENABLED=true +SHARE_ROLLUP_INTERVAL_MS=60000 +SHARE_ROLLUP_SAFETY_LAG_SECONDS=30 +SHARE_ROLLUP_MAX_SHARES_PER_BATCH=5000000 #redis REDIS_URL=redis://redis:6379 diff --git a/docker-compose.external-db.yml b/docker-compose.external-db.yml index 19bc11f..d91d88a 100644 --- a/docker-compose.external-db.yml +++ b/docker-compose.external-db.yml @@ -76,6 +76,10 @@ services: SHARE_ACCOUNTING_MAX_QUEUE_SIZE: ${SHARE_ACCOUNTING_MAX_QUEUE_SIZE:-50000} SHARE_ACCOUNTING_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MS:-2500} SHARE_ACCOUNTING_SUMMARY_CACHE_MAX: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MAX:-10000} + SHARE_ROLLUP_ENABLED: ${SHARE_ROLLUP_ENABLED:-true} + SHARE_ROLLUP_INTERVAL_MS: ${SHARE_ROLLUP_INTERVAL_MS:-60000} + SHARE_ROLLUP_SAFETY_LAG_SECONDS: ${SHARE_ROLLUP_SAFETY_LAG_SECONDS:-30} + SHARE_ROLLUP_MAX_SHARES_PER_BATCH: ${SHARE_ROLLUP_MAX_SHARES_PER_BATCH:-5000000} healthcheck: test: ["CMD-SHELL", "node -e \"const http=require('http'); const https=require('https'); const secure=process.env.API_SECURE==='true'; const client=secure?https:http; const req=client.get({hostname:'127.0.0.1',port:process.env.API_PORT||3334,path:'/api/network',rejectUnauthorized:false},res=>process.exit(res.statusCode<500?0:1)); req.on('error',()=>process.exit(1)); req.setTimeout(5000,()=>{req.destroy(); process.exit(1);});\""] interval: 30s diff --git a/full-setup/docker-compose-mainnet.yml b/full-setup/docker-compose-mainnet.yml index 440dc15..1d4037b 100644 --- a/full-setup/docker-compose-mainnet.yml +++ b/full-setup/docker-compose-mainnet.yml @@ -113,6 +113,10 @@ services: STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1} STRATUM_SOCKET_TIMEOUT_MS: ${STRATUM_SOCKET_TIMEOUT_MS:-3600000} STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS: ${STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS:-60000} + SHARE_ROLLUP_ENABLED: ${SHARE_ROLLUP_ENABLED:-true} + SHARE_ROLLUP_INTERVAL_MS: ${SHARE_ROLLUP_INTERVAL_MS:-60000} + SHARE_ROLLUP_SAFETY_LAG_SECONDS: ${SHARE_ROLLUP_SAFETY_LAG_SECONDS:-30} + SHARE_ROLLUP_MAX_SHARES_PER_BATCH: ${SHARE_ROLLUP_MAX_SHARES_PER_BATCH:-5000000} networks: bitcoin: diff --git a/src/ORM/_migrations/AcceptedShareRetentionCompression1780966200000.ts b/src/ORM/_migrations/AcceptedShareRetentionCompression1780966200000.ts new file mode 100644 index 0000000..63aef52 --- /dev/null +++ b/src/ORM/_migrations/AcceptedShareRetentionCompression1780966200000.ts @@ -0,0 +1,54 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +export class AcceptedShareRetentionCompression1780966200000 implements MigrationInterface { + public name = 'AcceptedShareRetentionCompression1780966200000'; + public transaction = false; + + public async up(queryRunner: QueryRunner): Promise { + await queryRunner.query(` + ALTER TABLE "accepted_share_entity" SET ( + timescaledb.compress, + timescaledb.compress_orderby = '"acceptedAt" DESC', + timescaledb.compress_segmentby = '"address","clientName"' + ) + `); + + await queryRunner.query(`SELECT remove_compression_policy('accepted_share_entity', if_exists => TRUE)`); + await queryRunner.query(` + SELECT add_compression_policy( + 'accepted_share_entity', + INTERVAL '24 hours', + if_not_exists => TRUE + ) + `); + + await queryRunner.query(`SELECT remove_retention_policy('accepted_share_entity', if_exists => TRUE)`); + await queryRunner.query(` + SELECT add_retention_policy( + 'accepted_share_entity', + INTERVAL '7 days', + if_not_exists => TRUE + ) + `); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.query(`SELECT remove_compression_policy('accepted_share_entity', if_exists => TRUE)`); + await queryRunner.query(` + SELECT add_compression_policy( + 'accepted_share_entity', + INTERVAL '1 day', + if_not_exists => TRUE + ) + `); + + await queryRunner.query(`SELECT remove_retention_policy('accepted_share_entity', if_exists => TRUE)`); + await queryRunner.query(` + SELECT add_retention_policy( + 'accepted_share_entity', + INTERVAL '30 days', + if_not_exists => TRUE + ) + `); + } +} diff --git a/src/ORM/_migrations/ShareRollupBatches1780962600000.ts b/src/ORM/_migrations/ShareRollupBatches1780962600000.ts new file mode 100644 index 0000000..79485d5 --- /dev/null +++ b/src/ORM/_migrations/ShareRollupBatches1780962600000.ts @@ -0,0 +1,62 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +export class ShareRollupBatches1780962600000 implements MigrationInterface { + public name = 'ShareRollupBatches1780962600000'; + public transaction = false; + + public async up(queryRunner: QueryRunner): Promise { + await queryRunner.query(` + CREATE TABLE IF NOT EXISTS "share_rollup_batch" ( + "id" bigserial PRIMARY KEY, + "startShareIndex" bigint NOT NULL, + "endShareIndex" bigint NOT NULL, + "startAcceptedAt" timestamptz NOT NULL, + "endAcceptedAt" timestamptz NOT NULL, + "acceptedShareCount" bigint NOT NULL, + "creditedDifficulty" numeric NOT NULL, + "status" varchar(16) NOT NULL DEFAULT 'finalized', + "createdAt" timestamptz NOT NULL DEFAULT CURRENT_TIMESTAMP, + "finalizedAt" timestamptz, + CONSTRAINT "CHK_share_rollup_batch_index_order" + CHECK ("endShareIndex" >= "startShareIndex") + ) + `); + await queryRunner.query(` + CREATE UNIQUE INDEX IF NOT EXISTS "IDX_share_rollup_batch_range" + ON "share_rollup_batch" ("startShareIndex", "endShareIndex") + `); + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS "IDX_share_rollup_batch_finalized_end" + ON "share_rollup_batch" ("endShareIndex" DESC) + WHERE "status" = 'finalized' + `); + + await queryRunner.query(` + CREATE TABLE IF NOT EXISTS "share_rollup_batch_summary" ( + "batchId" bigint NOT NULL REFERENCES "share_rollup_batch" ("id") ON DELETE CASCADE, + "address" varchar(62) NOT NULL, + "clientName" varchar NOT NULL, + "protocol" varchar(8) NOT NULL, + "blockHeight" integer NOT NULL, + "creditedDifficulty" numeric NOT NULL, + "acceptedShareCount" bigint NOT NULL, + "bestSubmissionDifficulty" numeric NOT NULL, + "firstShareAt" timestamptz NOT NULL, + "lastShareAt" timestamptz NOT NULL, + PRIMARY KEY ("batchId", "address", "clientName", "protocol", "blockHeight") + ) + `); + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS "IDX_share_rollup_summary_address_batch" + ON "share_rollup_batch_summary" ("address", "clientName", "blockHeight", "batchId") + `); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.query(`DROP INDEX IF EXISTS "IDX_share_rollup_summary_address_batch"`); + await queryRunner.query(`DROP TABLE IF EXISTS "share_rollup_batch_summary"`); + await queryRunner.query(`DROP INDEX IF EXISTS "IDX_share_rollup_batch_finalized_end"`); + await queryRunner.query(`DROP INDEX IF EXISTS "IDX_share_rollup_batch_range"`); + await queryRunner.query(`DROP TABLE IF EXISTS "share_rollup_batch"`); + } +} diff --git a/src/ORM/share-accounting/share-accounting.service.spec.ts b/src/ORM/share-accounting/share-accounting.service.spec.ts index 0b31d36..f19a65d 100644 --- a/src/ORM/share-accounting/share-accounting.service.spec.ts +++ b/src/ORM/share-accounting/share-accounting.service.spec.ts @@ -325,6 +325,69 @@ describe('ShareAccountingService', () => { expect(repository.query).toHaveBeenCalledTimes(1); }); + + it('should finalize accepted shares into an append-only share rollup batch', async () => { + const manager = { + query: jest.fn() + .mockResolvedValueOnce([{ locked: true }]) + .mockResolvedValueOnce([{ lastProcessedShareIndex: '10' }]) + .mockResolvedValueOnce([{ startShareIndex: '11', endShareIndex: '20', acceptedShareCount: 10 }]) + .mockResolvedValueOnce([{ + startAcceptedAt: new Date('2026-06-07T12:00:00Z'), + endAcceptedAt: new Date('2026-06-07T12:00:30Z'), + acceptedShareCount: '10', + creditedDifficulty: '2048', + }]) + .mockResolvedValueOnce([{ batchId: '7' }]) + .mockResolvedValueOnce([]), + }; + const repository = { + manager: { + transaction: jest.fn(callback => callback(manager)), + }, + }; + const service = new ShareAccountingService(repository as any); + + await expect(service.processPendingShareRollupBatch()).resolves.toEqual({ + processed: true, + batchId: '7', + startShareIndex: '11', + endShareIndex: '20', + acceptedShareCount: 10, + creditedDifficulty: 2048, + }); + + expect(repository.manager.transaction).toHaveBeenCalledTimes(1); + expect(manager.query).toHaveBeenNthCalledWith( + 1, + expect.stringContaining('pg_try_advisory_xact_lock'), + ['1780962600'], + ); + expect(manager.query).toHaveBeenNthCalledWith( + 6, + expect.stringContaining('INSERT INTO "share_rollup_batch_summary"'), + ['7', '11', '20'], + ); + }); + + it('should skip share rollup work when another process holds the advisory lock', async () => { + const manager = { + query: jest.fn().mockResolvedValueOnce([{ locked: false }]), + }; + const repository = { + manager: { + transaction: jest.fn(callback => callback(manager)), + }, + }; + const service = new ShareAccountingService(repository as any); + + await expect(service.processPendingShareRollupBatch()).resolves.toEqual({ + processed: false, + reason: 'locked', + }); + + expect(manager.query).toHaveBeenCalledTimes(1); + }); }); function buildRecord(jobId: string) { diff --git a/src/ORM/share-accounting/share-accounting.service.ts b/src/ORM/share-accounting/share-accounting.service.ts index 5bdde7a..d3c836a 100644 --- a/src/ORM/share-accounting/share-accounting.service.ts +++ b/src/ORM/share-accounting/share-accounting.service.ts @@ -1,4 +1,4 @@ -import { Injectable, OnModuleDestroy } from '@nestjs/common'; +import { Injectable, OnModuleDestroy, OnModuleInit } from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; import { Repository } from 'typeorm'; @@ -59,6 +59,16 @@ export interface SessionShareSummary { bestSubmissionDifficulty: number; } +export interface ShareRollupBatchResult { + processed: boolean; + reason?: 'disabled' | 'locked' | 'no-shares'; + batchId?: string; + startShareIndex?: string; + endShareIndex?: string; + acceptedShareCount?: number; + creditedDifficulty?: number; +} + interface AccountingFilter { address?: string; clientName?: string; @@ -73,12 +83,18 @@ 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; +const DEFAULT_ROLLUP_INTERVAL_MS = 60000; +const DEFAULT_ROLLUP_SAFETY_LAG_SECONDS = 10; +const DEFAULT_ROLLUP_MAX_SHARES_PER_BATCH = 5000000; +const SHARE_ROLLUP_ADVISORY_LOCK = '1780962600'; @Injectable() -export class ShareAccountingService implements OnModuleDestroy { +export class ShareAccountingService implements OnModuleInit, OnModuleDestroy { private pendingShares: PendingShare[] = []; private flushTimer: NodeJS.Timeout | null = null; private activeFlush: Promise | null = null; + private rollupTimer: NodeJS.Timeout | null = null; + private activeRollup: Promise | null = null; private summaryCache = new Map(); private readonly poolSummaryCacheKey = 'accounting:pool-summary'; private readonly batchSize = this.readPositiveInt('SHARE_ACCOUNTING_BATCH_SIZE', DEFAULT_BATCH_SIZE); @@ -87,6 +103,10 @@ export class ShareAccountingService implements OnModuleDestroy { 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); + private readonly shareRollupEnabled = this.readBoolean('SHARE_ROLLUP_ENABLED', true); + private readonly shareRollupIntervalMs = this.readPositiveInt('SHARE_ROLLUP_INTERVAL_MS', DEFAULT_ROLLUP_INTERVAL_MS); + private readonly shareRollupSafetyLagSeconds = this.readNonNegativeInt('SHARE_ROLLUP_SAFETY_LAG_SECONDS', DEFAULT_ROLLUP_SAFETY_LAG_SECONDS); + private readonly shareRollupMaxSharesPerBatch = this.readPositiveInt('SHARE_ROLLUP_MAX_SHARES_PER_BATCH', DEFAULT_ROLLUP_MAX_SHARES_PER_BATCH); constructor( @InjectRepository(AcceptedShareEntity) @@ -94,6 +114,10 @@ export class ShareAccountingService implements OnModuleDestroy { private readonly redisMessagingService?: RedisMessagingService, ) { } + public onModuleInit(): void { + this.startShareRollupTimer(); + } + public async recordAcceptedShare(record: AcceptedShareRecord): Promise { const acceptedShare = this.acceptedShareRepository.create({ ...record, @@ -134,9 +158,171 @@ export class ShareAccountingService implements OnModuleDestroy { } public async onModuleDestroy(): Promise { + if (this.rollupTimer != null) { + clearInterval(this.rollupTimer); + this.rollupTimer = null; + } + if (this.activeRollup != null) { + await this.activeRollup; + } await this.flushPendingShares(); } + public async processPendingShareRollupBatch(): Promise { + if (!this.shareRollupEnabled) { + return { processed: false, reason: 'disabled' }; + } + + return this.acceptedShareRepository.manager.transaction(async manager => { + const [lockRow] = await manager.query(` + SELECT pg_try_advisory_xact_lock($1::bigint) AS "locked" + `, [SHARE_ROLLUP_ADVISORY_LOCK]); + + if (lockRow?.locked !== true) { + return { processed: false, reason: 'locked' }; + } + + const [lastRow] = await manager.query(` + SELECT COALESCE(MAX("endShareIndex"), 0)::text AS "lastProcessedShareIndex" + FROM "share_rollup_batch" + WHERE "status" = 'finalized' + `); + const lastProcessedShareIndex = lastRow?.lastProcessedShareIndex ?? '0'; + + const [rangeRow] = await manager.query(` + WITH last_state AS MATERIALIZED ( + SELECT $1::bigint AS "lastProcessedShareIndex" + ), + first_unstable AS MATERIALIZED ( + SELECT MIN("shareIndex") AS "firstUnstableShareIndex" + FROM "accepted_share_entity", last_state + WHERE "shareIndex" > last_state."lastProcessedShareIndex" + AND "acceptedAt" > NOW() - ($2::int * INTERVAL '1 second') + ), + stable_bound AS MATERIALIZED ( + SELECT + CASE + WHEN (SELECT "firstUnstableShareIndex" FROM first_unstable) IS NULL THEN ( + SELECT MAX("shareIndex") + FROM "accepted_share_entity", last_state + WHERE "shareIndex" > last_state."lastProcessedShareIndex" + ) + ELSE (SELECT "firstUnstableShareIndex" FROM first_unstable) - 1 + END AS "maxStableShareIndex" + ), + bounded AS MATERIALIZED ( + SELECT "shareIndex" + FROM "accepted_share_entity", last_state, stable_bound + WHERE "shareIndex" > last_state."lastProcessedShareIndex" + AND "shareIndex" <= stable_bound."maxStableShareIndex" + ORDER BY "shareIndex" ASC + LIMIT $3::int + ) + SELECT + MIN("shareIndex")::text AS "startShareIndex", + MAX("shareIndex")::text AS "endShareIndex", + COUNT(*)::int AS "acceptedShareCount" + FROM bounded + `, [ + lastProcessedShareIndex, + this.shareRollupSafetyLagSeconds, + this.shareRollupMaxSharesPerBatch, + ]); + + if (rangeRow?.startShareIndex == null || rangeRow?.endShareIndex == null || this.toNumber(rangeRow.acceptedShareCount) === 0) { + return { processed: false, reason: 'no-shares' }; + } + + const [statsRow] = await manager.query(` + SELECT + MIN("acceptedAt") AS "startAcceptedAt", + MAX("acceptedAt") AS "endAcceptedAt", + COUNT(*)::bigint AS "acceptedShareCount", + COALESCE(SUM("creditedDifficulty"), 0)::numeric AS "creditedDifficulty" + FROM "accepted_share_entity" + WHERE "shareIndex" >= $1::bigint + AND "shareIndex" <= $2::bigint + `, [rangeRow.startShareIndex, rangeRow.endShareIndex]); + + if (statsRow?.startAcceptedAt == null || this.toNumber(statsRow.acceptedShareCount) === 0) { + return { processed: false, reason: 'no-shares' }; + } + + const [batchRow] = await manager.query(` + INSERT INTO "share_rollup_batch" ( + "startShareIndex", + "endShareIndex", + "startAcceptedAt", + "endAcceptedAt", + "acceptedShareCount", + "creditedDifficulty", + "status", + "finalizedAt" + ) VALUES ( + $1::bigint, + $2::bigint, + $3::timestamptz, + $4::timestamptz, + $5::bigint, + $6::numeric, + 'finalized', + NOW() + ) + RETURNING "id"::text AS "batchId" + `, [ + rangeRow.startShareIndex, + rangeRow.endShareIndex, + statsRow.startAcceptedAt, + statsRow.endAcceptedAt, + statsRow.acceptedShareCount, + statsRow.creditedDifficulty, + ]); + + await manager.query(` + INSERT INTO "share_rollup_batch_summary" ( + "batchId", + "address", + "clientName", + "protocol", + "blockHeight", + "creditedDifficulty", + "acceptedShareCount", + "bestSubmissionDifficulty", + "firstShareAt", + "lastShareAt" + ) + SELECT + $1::bigint AS "batchId", + "address", + "clientName", + "protocol", + "blockHeight", + SUM("creditedDifficulty") AS "creditedDifficulty", + COUNT(*)::bigint AS "acceptedShareCount", + MAX("submissionDifficulty") AS "bestSubmissionDifficulty", + MIN("acceptedAt") AS "firstShareAt", + MAX("acceptedAt") AS "lastShareAt" + FROM "accepted_share_entity" + WHERE "shareIndex" >= $2::bigint + AND "shareIndex" <= $3::bigint + GROUP BY "address", "clientName", "protocol", "blockHeight" + `, [ + batchRow.batchId, + rangeRow.startShareIndex, + rangeRow.endShareIndex, + ]); + + return { + processed: true, + batchId: batchRow.batchId, + startShareIndex: rangeRow.startShareIndex, + endShareIndex: rangeRow.endShareIndex, + acceptedShareCount: this.toNumber(statsRow.acceptedShareCount), + creditedDifficulty: this.toNumber(statsRow.creditedDifficulty), + }; + }); + } + public async getPoolSummary(): Promise { const cached = await this.redisMessagingService ?.getJsonCache(this.poolSummaryCacheKey) @@ -473,6 +659,32 @@ export class ShareAccountingService implements OnModuleDestroy { } } + private startShareRollupTimer(): void { + if (!this.shareRollupEnabled || process.env.MASTER !== 'true') { + return; + } + + this.rollupTimer = setInterval(() => { + if (this.activeRollup != null) { + return; + } + + this.activeRollup = this.processPendingShareRollupBatch() + .then(result => { + if (result.processed) { + console.log(`Share rollup batch ${result.batchId} finalized: indexes ${result.startShareIndex}-${result.endShareIndex}, shares ${result.acceptedShareCount}, difficulty ${result.creditedDifficulty}`); + } + }) + .catch(error => { + console.error(`Share rollup batch failed: ${error.message}`); + }) + .finally(() => { + this.activeRollup = null; + }); + }, this.shareRollupIntervalMs); + this.rollupTimer.unref?.(); + } + private buildWhereClause(filter: AccountingFilter): { whereSql: string; params: string[] } { const where: string[] = []; const params: string[] = []; @@ -532,6 +744,14 @@ export class ShareAccountingService implements OnModuleDestroy { const value = Number(process.env[name]); return Number.isInteger(value) && value >= 0 ? value : defaultValue; } + + private readBoolean(name: string, defaultValue: boolean): boolean { + const value = process.env[name]?.toLowerCase(); + if (value == null || value.length === 0) { + return defaultValue; + } + return value === 'true' || value === '1' || value === 'yes'; + } } interface PendingShare { diff --git a/src/database.config.ts b/src/database.config.ts index 871139e..fd22222 100644 --- a/src/database.config.ts +++ b/src/database.config.ts @@ -9,6 +9,8 @@ import { AcceptedShareRollupIndexes1780865400000 } from './ORM/_migrations/Accep import { PoolAccountingDashboardIndexes1780867200000 } from './ORM/_migrations/PoolAccountingDashboardIndexes1780867200000'; import { CurrentRoundBestShareIndex1780897600000 } from './ORM/_migrations/CurrentRoundBestShareIndex1780897600000'; import { CurrentRoundWorkRollup1780899000000 } from './ORM/_migrations/CurrentRoundWorkRollup1780899000000'; +import { ShareRollupBatches1780962600000 } from './ORM/_migrations/ShareRollupBatches1780962600000'; +import { AcceptedShareRetentionCompression1780966200000 } from './ORM/_migrations/AcceptedShareRetentionCompression1780966200000'; 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'; @@ -36,6 +38,8 @@ export const databaseMigrations = [ PoolAccountingDashboardIndexes1780867200000, CurrentRoundBestShareIndex1780897600000, CurrentRoundWorkRollup1780899000000, + ShareRollupBatches1780962600000, + AcceptedShareRetentionCompression1780966200000, ]; 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 3ae3b86..d5227a2 100644 --- a/test/timescale-redis.integration-spec.ts +++ b/test/timescale-redis.integration-spec.ts @@ -42,6 +42,8 @@ describe('TimescaleDB and Redis integration', () => { }); beforeEach(async () => { + await dataSource.query(`DELETE FROM share_rollup_batch_summary`); + await dataSource.query(`DELETE FROM share_rollup_batch`); await dataSource.query(`DELETE FROM accepted_share_entity`); await dataSource.query(`DELETE FROM client_entity`); await dataSource.query(`REFRESH MATERIALIZED VIEW user_agent_report_view`); @@ -103,6 +105,18 @@ describe('TimescaleDB and Redis integration', () => { sequence_name: 'accepted_share_index_seq', index_name: '"IDX_accepted_share_order"', }); + + const shareRollupObjects = await dataSource.query(` + SELECT + to_regclass('public.share_rollup_batch') AS batch_table, + to_regclass('public.share_rollup_batch_summary') AS summary_table, + to_regclass('public."IDX_share_rollup_batch_finalized_end"') AS finalized_index + `); + expect(shareRollupObjects[0]).toEqual({ + batch_table: 'share_rollup_batch', + summary_table: 'share_rollup_batch_summary', + finalized_index: '"IDX_share_rollup_batch_finalized_end"', + }); }); it('should persist accepted shares and refresh the 10 minute aggregate', async () => { @@ -192,6 +206,99 @@ describe('TimescaleDB and Redis integration', () => { ])); }); + it('should finalize accepted shares into share rollup batches without mutating raw rows', async () => { + const client = await dataSource.getRepository(ClientEntity).save({ + address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', + clientName: 'rollup-worker', + sessionId: '47a6f098', + userAgent: 'integration-test', + startTime: new Date(), + bestDifficulty: 0, + hashRate: 0, + }); + const service = new ShareAccountingService(dataSource.getRepository(AcceptedShareEntity)); + const acceptedAt = new Date('2026-06-07T12:20:00Z'); + + await service.recordAcceptedShare({ + protocol: 'sv1', + acceptedAt, + address: client.address, + clientName: client.clientName, + sessionId: client.sessionId, + clientId: client.id, + jobId: 'rollup-1', + jobTemplateId: 'rollup-template', + blockHeight: 900010, + creditedDifficulty: 64, + submissionDifficulty: 128, + networkDifficulty: 100000, + nonce: 'rollup-nonce-1', + ntime: '64b3f3ec', + version: '20000000', + extraNonce2: 'c708000000000001', + isBlockCandidate: false, + blockSubmissionResult: null, + }); + await service.recordAcceptedShare({ + protocol: 'sv2', + acceptedAt: new Date(acceptedAt.getTime() + 1), + address: client.address, + clientName: client.clientName, + sessionId: client.sessionId, + clientId: client.id, + jobId: 'rollup-2', + jobTemplateId: 'rollup-template', + blockHeight: 900010, + creditedDifficulty: 32, + submissionDifficulty: 64, + networkDifficulty: 100000, + nonce: 'rollup-nonce-2', + ntime: '64b3f3ec', + version: '20000000', + extraNonce2: 'c708000000000002', + isBlockCandidate: false, + blockSubmissionResult: null, + }); + + await expect(service.processPendingShareRollupBatch()).resolves.toEqual(expect.objectContaining({ + processed: true, + acceptedShareCount: 2, + creditedDifficulty: 96, + })); + await expect(service.processPendingShareRollupBatch()).resolves.toEqual({ + processed: false, + reason: 'no-shares', + }); + + const batches = await dataSource.query(` + SELECT "acceptedShareCount"::int AS "acceptedShareCount", "creditedDifficulty"::float AS "creditedDifficulty" + FROM share_rollup_batch + `); + expect(batches).toEqual([{ + acceptedShareCount: 2, + creditedDifficulty: 96, + }]); + + const summaries = await dataSource.query(` + SELECT + "protocol", + "blockHeight", + "acceptedShareCount"::int AS "acceptedShareCount", + "creditedDifficulty"::float AS "creditedDifficulty", + "bestSubmissionDifficulty"::float AS "bestSubmissionDifficulty" + FROM share_rollup_batch_summary + WHERE "address" = $1 AND "clientName" = $2 + ORDER BY "protocol" + `, [client.address, client.clientName]); + expect(summaries).toEqual([ + { protocol: 'sv1', blockHeight: 900010, acceptedShareCount: 1, creditedDifficulty: 64, bestSubmissionDifficulty: 128 }, + { protocol: 'sv2', blockHeight: 900010, acceptedShareCount: 1, creditedDifficulty: 32, bestSubmissionDifficulty: 64 }, + ]); + + const rawRows = await dataSource.query(`SELECT COUNT(*)::int AS count FROM accepted_share_entity`); + expect(rawRows[0].count).toBe(2); + }); + it('should exclude soft-deleted clients from the user-agent report', async () => { const userAgent = `integration-reconnect-${Date.now()}`; const repository = dataSource.getRepository(ClientEntity);