Add accepted share rollup batches

This commit is contained in:
Ben
2026-06-09 22:55:11 -04:00
parent 1c5a3d9a94
commit bcd6bb0aa7
9 changed files with 524 additions and 2 deletions
@@ -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<void> {
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<void> {
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
)
`);
}
}
@@ -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<void> {
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<void> {
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"`);
}
}
@@ -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) {
@@ -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<void> | null = null;
private rollupTimer: NodeJS.Timeout | null = null;
private activeRollup: Promise<void> | null = null;
private summaryCache = new Map<string, SummaryCacheEntry>();
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<AcceptedShareEntity> {
const acceptedShare = this.acceptedShareRepository.create({
...record,
@@ -134,9 +158,171 @@ export class ShareAccountingService implements OnModuleDestroy {
}
public async onModuleDestroy(): Promise<void> {
if (this.rollupTimer != null) {
clearInterval(this.rollupTimer);
this.rollupTimer = null;
}
if (this.activeRollup != null) {
await this.activeRollup;
}
await this.flushPendingShares();
}
public async processPendingShareRollupBatch(): Promise<ShareRollupBatchResult> {
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<ShareAccountingSummary> {
const cached = await this.redisMessagingService
?.getJsonCache<ShareAccountingSummary>(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 {
+4
View File
@@ -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 {