Compare commits

2 Commits
Author SHA1 Message Date
Ben bcd6bb0aa7 Add accepted share rollup batches 2026-06-09 22:55:11 -04:00
Ben 1c5a3d9a94 performance and optimization 2026-06-09 17:13:27 -04:00
22 changed files with 808 additions and 50 deletions
+4
View File
@@ -77,6 +77,10 @@ SHARE_ACCOUNTING_FLUSH_INTERVAL_MS=25
SHARE_ACCOUNTING_MAX_QUEUE_SIZE=50000 SHARE_ACCOUNTING_MAX_QUEUE_SIZE=50000
SHARE_ACCOUNTING_SUMMARY_CACHE_MS=2500 SHARE_ACCOUNTING_SUMMARY_CACHE_MS=2500
SHARE_ACCOUNTING_SUMMARY_CACHE_MAX=10000 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
REDIS_URL=redis://redis:6379 REDIS_URL=redis://redis:6379
+1
View File
@@ -20,6 +20,7 @@ lerna-debug.log*
/.nyc_output /.nyc_output
/.baseline-*.json /.baseline-*.json
/.test-artifacts /.test-artifacts
/.profiles
# IDEs and editors # IDEs and editors
/.idea /.idea
+4
View File
@@ -76,6 +76,10 @@ services:
SHARE_ACCOUNTING_MAX_QUEUE_SIZE: ${SHARE_ACCOUNTING_MAX_QUEUE_SIZE:-50000} 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_MS: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MS:-2500}
SHARE_ACCOUNTING_SUMMARY_CACHE_MAX: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MAX:-10000} 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: 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);});\""] 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 interval: 30s
+4
View File
@@ -113,6 +113,10 @@ services:
STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1} STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1}
STRATUM_SOCKET_TIMEOUT_MS: ${STRATUM_SOCKET_TIMEOUT_MS:-3600000} STRATUM_SOCKET_TIMEOUT_MS: ${STRATUM_SOCKET_TIMEOUT_MS:-3600000}
STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS: ${STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS:-60000} 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: networks:
bitcoin: bitcoin:
@@ -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); 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) { 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 { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm'; import { Repository } from 'typeorm';
@@ -59,6 +59,16 @@ export interface SessionShareSummary {
bestSubmissionDifficulty: number; bestSubmissionDifficulty: number;
} }
export interface ShareRollupBatchResult {
processed: boolean;
reason?: 'disabled' | 'locked' | 'no-shares';
batchId?: string;
startShareIndex?: string;
endShareIndex?: string;
acceptedShareCount?: number;
creditedDifficulty?: number;
}
interface AccountingFilter { interface AccountingFilter {
address?: string; address?: string;
clientName?: string; clientName?: string;
@@ -73,12 +83,18 @@ const DEFAULT_MAX_QUEUE_SIZE = 50000;
const DEFAULT_SUMMARY_CACHE_MS = 2500; const DEFAULT_SUMMARY_CACHE_MS = 2500;
const DEFAULT_SUMMARY_CACHE_MAX = 10000; const DEFAULT_SUMMARY_CACHE_MAX = 10000;
const DEFAULT_REDIS_SUMMARY_CACHE_MS = 30000; 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() @Injectable()
export class ShareAccountingService implements OnModuleDestroy { export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
private pendingShares: PendingShare[] = []; private pendingShares: PendingShare[] = [];
private flushTimer: NodeJS.Timeout | null = null; private flushTimer: NodeJS.Timeout | null = null;
private activeFlush: Promise<void> | 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 summaryCache = new Map<string, SummaryCacheEntry>();
private readonly poolSummaryCacheKey = 'accounting:pool-summary'; private readonly poolSummaryCacheKey = 'accounting:pool-summary';
private readonly batchSize = this.readPositiveInt('SHARE_ACCOUNTING_BATCH_SIZE', DEFAULT_BATCH_SIZE); 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 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 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 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( constructor(
@InjectRepository(AcceptedShareEntity) @InjectRepository(AcceptedShareEntity)
@@ -94,6 +114,10 @@ export class ShareAccountingService implements OnModuleDestroy {
private readonly redisMessagingService?: RedisMessagingService, private readonly redisMessagingService?: RedisMessagingService,
) { } ) { }
public onModuleInit(): void {
this.startShareRollupTimer();
}
public async recordAcceptedShare(record: AcceptedShareRecord): Promise<AcceptedShareEntity> { public async recordAcceptedShare(record: AcceptedShareRecord): Promise<AcceptedShareEntity> {
const acceptedShare = this.acceptedShareRepository.create({ const acceptedShare = this.acceptedShareRepository.create({
...record, ...record,
@@ -134,9 +158,171 @@ export class ShareAccountingService implements OnModuleDestroy {
} }
public async onModuleDestroy(): Promise<void> { 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(); 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> { public async getPoolSummary(): Promise<ShareAccountingSummary> {
const cached = await this.redisMessagingService const cached = await this.redisMessagingService
?.getJsonCache<ShareAccountingSummary>(this.poolSummaryCacheKey) ?.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[] } { private buildWhereClause(filter: AccountingFilter): { whereSql: string; params: string[] } {
const where: string[] = []; const where: string[] = [];
const params: string[] = []; const params: string[] = [];
@@ -532,6 +744,14 @@ export class ShareAccountingService implements OnModuleDestroy {
const value = Number(process.env[name]); const value = Number(process.env[name]);
return Number.isInteger(value) && value >= 0 ? value : defaultValue; 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 { interface PendingShare {
+4
View File
@@ -9,6 +9,8 @@ import { AcceptedShareRollupIndexes1780865400000 } from './ORM/_migrations/Accep
import { PoolAccountingDashboardIndexes1780867200000 } from './ORM/_migrations/PoolAccountingDashboardIndexes1780867200000'; import { PoolAccountingDashboardIndexes1780867200000 } from './ORM/_migrations/PoolAccountingDashboardIndexes1780867200000';
import { CurrentRoundBestShareIndex1780897600000 } from './ORM/_migrations/CurrentRoundBestShareIndex1780897600000'; import { CurrentRoundBestShareIndex1780897600000 } from './ORM/_migrations/CurrentRoundBestShareIndex1780897600000';
import { CurrentRoundWorkRollup1780899000000 } from './ORM/_migrations/CurrentRoundWorkRollup1780899000000'; 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 { 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';
@@ -36,6 +38,8 @@ export const databaseMigrations = [
PoolAccountingDashboardIndexes1780867200000, PoolAccountingDashboardIndexes1780867200000,
CurrentRoundBestShareIndex1780897600000, CurrentRoundBestShareIndex1780897600000,
CurrentRoundWorkRollup1780899000000, CurrentRoundWorkRollup1780899000000,
ShareRollupBatches1780962600000,
AcceptedShareRetentionCompression1780966200000,
]; ];
export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions { export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions {
+3 -2
View File
@@ -2,6 +2,7 @@ import { AddressType, getAddressInfo } from 'bitcoin-address-validation';
import * as bitcoinjs from 'bitcoinjs-lib'; import * as bitcoinjs from 'bitcoinjs-lib';
import { IJobTemplate } from '../services/stratum-v1-jobs.service'; import { IJobTemplate } from '../services/stratum-v1-jobs.service';
import { hash256 } from '../utils/hash.utils';
import { eResponseMethod } from './enums/eResponseMethod'; import { eResponseMethod } from './enums/eResponseMethod';
import { IMiningNotify } from './stratum-messages/IMiningNotify'; import { IMiningNotify } from './stratum-messages/IMiningNotify';
import { TOTAL_EXTRANONCE_SIZE_BYTES } from './stratum.constants'; import { TOTAL_EXTRANONCE_SIZE_BYTES } from './stratum.constants';
@@ -96,7 +97,7 @@ export class MiningJob {
Buffer.from(`${extraNonce}${extraNonce2}`, 'hex'), Buffer.from(`${extraNonce}${extraNonce2}`, 'hex'),
this.coinbasePart2Buffer, this.coinbasePart2Buffer,
]); ]);
const coinbaseHash = bitcoinjs.crypto.hash256(coinbaseBuffer); const coinbaseHash = hash256(coinbaseBuffer);
const merkleRoot = this.calculateMerkleRootHash(coinbaseHash, this.merkleBranchBuffers); const merkleRoot = this.calculateMerkleRootHash(coinbaseHash, this.merkleBranchBuffers);
let version = jobTemplate.block.version; let version = jobTemplate.block.version;
@@ -154,7 +155,7 @@ export class MiningJob {
for (let i = 0; i < merkleBranches.length; i++) { for (let i = 0; i < merkleBranches.length; i++) {
bothMerkles.set(merkleBranches[i], 32); bothMerkles.set(merkleBranches[i], 32);
newRoot = bitcoinjs.crypto.hash256(bothMerkles); newRoot = hash256(bothMerkles);
bothMerkles.set(newRoot); bothMerkles.set(newRoot);
} }
+26 -2
View File
@@ -9,6 +9,7 @@ import { ClientService } from '../ORM/client/client.service';
import { BitcoinRpcService as MockBitcoinRpcService } from '../services/bitcoin-rpc.service'; import { BitcoinRpcService as MockBitcoinRpcService } from '../services/bitcoin-rpc.service';
import { NotificationService } from '../services/notification.service'; import { NotificationService } from '../services/notification.service';
import { StratumV1JobsService } from '../services/stratum-v1-jobs.service'; import { StratumV1JobsService } from '../services/stratum-v1-jobs.service';
import { DifficultyUtils } from '../utils/difficulty.utils';
import { IBlockTemplate } from './bitcoin-rpc/IBlockTemplate'; import { IBlockTemplate } from './bitcoin-rpc/IBlockTemplate';
import { MiningJob } from './MiningJob'; import { MiningJob } from './MiningJob';
import { StratumV1Client } from './StratumV1Client'; import { StratumV1Client } from './StratumV1Client';
@@ -471,7 +472,8 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
jest.spyOn(client as any, 'calculateDifficulty').mockReturnValue({ jest.spyOn(client as any, 'calculateDifficulty').mockReturnValue({
submissionDifficulty: 1024, submissionDifficulty: 1024,
submissionHash: 'share' submissionHash: 'share',
hashBuffer: DifficultyUtils.difficultyToTarget(1024),
}); });
const getSettingsSpy = jest.spyOn(addressSettings, 'getSettings'); const getSettingsSpy = jest.spyOn(addressSettings, 'getSettings');
const updateIfHigherSpy = jest.spyOn(addressSettings as any, 'updateBestDifficultyIfHigher').mockResolvedValue({ affected: 1 }); const updateIfHigherSpy = jest.spyOn(addressSettings as any, 'updateBestDifficultyIfHigher').mockResolvedValue({ affected: 1 });
@@ -495,6 +497,27 @@ describe('StratumV1Client', () => {
expect(getSettingsSpy).not.toHaveBeenCalled(); expect(getSettingsSpy).not.toHaveBeenCalled();
}); });
it('should reject shares by exact target even when reported difficulty is huge', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
jest.spyOn(client as any, 'calculateDifficulty').mockReturnValue({
submissionDifficulty: Number.MAX_SAFE_INTEGER,
submissionHash: 'too-easy',
hashBuffer: Buffer.alloc(32, 0xff),
});
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1024]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
emitMessage(MockRecording1.MINING_SUBMIT);
jest.useRealTimers();
await new Promise((r) => setTimeout(r, 1000));
expect((client as any).write).lastCalledWith(`{"id":5,"result":null,"error":[23,"Difficulty too low",""]}\n`);
expect(shareAccountingService.recordAcceptedShare).not.toHaveBeenCalled();
});
it('should reject duplicate submissions', async () => { it('should reject duplicate submissions', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
@@ -626,7 +649,8 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
jest.spyOn(client as any, 'calculateDifficulty').mockReturnValue({ jest.spyOn(client as any, 'calculateDifficulty').mockReturnValue({
submissionDifficulty: Number.MAX_SAFE_INTEGER, submissionDifficulty: Number.MAX_SAFE_INTEGER,
submissionHash: 'block-share' submissionHash: 'block-share',
hashBuffer: Buffer.alloc(32),
}); });
jest.spyOn(addressSettings, 'resetBestDifficultyAndShares').mockResolvedValue(undefined); jest.spyOn(addressSettings, 'resetBestDifficultyAndShares').mockResolvedValue(undefined);
+17 -7
View File
@@ -16,6 +16,8 @@ import { BitcoinRpcService } from '../services/bitcoin-rpc.service';
import { NotificationService } from '../services/notification.service'; import { NotificationService } from '../services/notification.service';
import { RedisMessagingService } from '../services/redis-messaging.service'; import { RedisMessagingService } from '../services/redis-messaging.service';
import { IJobTemplate, StratumV1JobsService } from '../services/stratum-v1-jobs.service'; import { IJobTemplate, StratumV1JobsService } from '../services/stratum-v1-jobs.service';
import { DifficultyUtils } from '../utils/difficulty.utils';
import { hash256 } from '../utils/hash.utils';
import { eRequestMethod } from './enums/eRequestMethod'; import { eRequestMethod } from './enums/eRequestMethod';
import { eResponseMethod } from './enums/eResponseMethod'; import { eResponseMethod } from './enums/eResponseMethod';
import { eStratumErrorCode } from './enums/eStratumErrorCode'; import { eStratumErrorCode } from './enums/eStratumErrorCode';
@@ -51,6 +53,7 @@ export class StratumV1Client {
private stratumInitialized = false; private stratumInitialized = false;
private usedSuggestedDifficulty = false; private usedSuggestedDifficulty = false;
private sessionDifficulty: number = 100000; private sessionDifficulty: number = 100000;
private sessionDifficultyTarget: Buffer = DifficultyUtils.difficultyToTarget(this.sessionDifficulty);
private clientEntity: ClientEntity; private clientEntity: ClientEntity;
private creatingEntity: Promise<void>; private creatingEntity: Promise<void>;
@@ -268,6 +271,7 @@ export class StratumV1Client {
this.clientAuthorization = authorizationMessage; this.clientAuthorization = authorizationMessage;
if (this.clientSuggestedDifficulty == null && this.clientAuthorization.startingDiff != null && this.clientAuthorization.startingDiff > this.sessionDifficulty) { if (this.clientSuggestedDifficulty == null && this.clientAuthorization.startingDiff != null && this.clientAuthorization.startingDiff > this.sessionDifficulty) {
this.sessionDifficulty = this.clientAuthorization.startingDiff; this.sessionDifficulty = this.clientAuthorization.startingDiff;
this.sessionDifficultyTarget = DifficultyUtils.difficultyToTarget(this.sessionDifficulty);
} }
const success = await this.write(JSON.stringify(this.clientAuthorization.response()) + '\n'); const success = await this.write(JSON.stringify(this.clientAuthorization.response()) + '\n');
if (!success) { if (!success) {
@@ -309,6 +313,7 @@ export class StratumV1Client {
this.clientSuggestedDifficulty = suggestDifficultyMessage; this.clientSuggestedDifficulty = suggestDifficultyMessage;
this.sessionDifficulty = this.clampDifficulty(suggestDifficultyMessage.suggestedDifficulty); this.sessionDifficulty = this.clampDifficulty(suggestDifficultyMessage.suggestedDifficulty);
this.sessionDifficultyTarget = DifficultyUtils.difficultyToTarget(this.sessionDifficulty);
const success = await this.write(JSON.stringify(this.clientSuggestedDifficulty.response(this.sessionDifficulty)) + '\n'); const success = await this.write(JSON.stringify(this.clientSuggestedDifficulty.response(this.sessionDifficulty)) + '\n');
if (!success) { if (!success) {
return; return;
@@ -610,19 +615,23 @@ export class StratumV1Client {
submission.extraNonce2, submission.extraNonce2,
timestamp timestamp
); );
const { submissionDifficulty } = this.calculateDifficulty(header); const { submissionDifficulty, hashBuffer } = this.calculateDifficulty(header);
//console.log(`DIFF: ${submissionDifficulty} of ${this.sessionDifficulty} from ${this.clientAuthorization.worker + '.' + this.extraNonceAndSessionId}`); //console.log(`DIFF: ${submissionDifficulty} of ${this.sessionDifficulty} from ${this.clientAuthorization.worker + '.' + this.extraNonceAndSessionId}`);
if (submissionDifficulty >= this.sessionDifficulty) { if (DifficultyUtils.meetsTarget(hashBuffer, this.sessionDifficultyTarget)) {
const success = await this.write(JSON.stringify(submission.response()) + '\n'); const success = await this.write(JSON.stringify(submission.response()) + '\n');
if (!success) { if (!success) {
return false; return false;
} }
let blockSubmissionResult: string = null; let blockSubmissionResult: string = null;
if (submissionDifficulty >= jobTemplate.blockData.networkDifficulty) { const isBlockCandidate = DifficultyUtils.meetsTarget(
hashBuffer,
DifficultyUtils.difficultyToTarget(jobTemplate.blockData.networkDifficulty),
);
if (isBlockCandidate) {
console.log('!!! BLOCK FOUND !!!'); console.log('!!! BLOCK FOUND !!!');
const updatedJobBlock = job.copyAndUpdateBlock( const updatedJobBlock = job.copyAndUpdateBlock(
jobTemplate, jobTemplate,
@@ -668,7 +677,7 @@ export class StratumV1Client {
? (jobTemplate.block.version ^ versionMask).toString(16) ? (jobTemplate.block.version ^ versionMask).toString(16)
: jobTemplate.block.version.toString(16), : jobTemplate.block.version.toString(16),
extraNonce2: submission.extraNonce2, extraNonce2: submission.extraNonce2,
isBlockCandidate: submissionDifficulty >= jobTemplate.blockData.networkDifficulty, isBlockCandidate,
blockSubmissionResult, blockSubmissionResult,
}); });
await this.statistics.addShares(this.clientEntity, this.sessionDifficulty); await this.statistics.addShares(this.clientEntity, this.sessionDifficulty);
@@ -718,6 +727,7 @@ export class StratumV1Client {
if (targetDiff != this.sessionDifficulty) { if (targetDiff != this.sessionDifficulty) {
//console.log(`Adjusting ${this.extraNonceAndSessionId} difficulty from ${this.sessionDifficulty} to ${targetDiff}`); //console.log(`Adjusting ${this.extraNonceAndSessionId} difficulty from ${this.sessionDifficulty} to ${targetDiff}`);
this.sessionDifficulty = targetDiff; this.sessionDifficulty = targetDiff;
this.sessionDifficultyTarget = DifficultyUtils.difficultyToTarget(this.sessionDifficulty);
const data = JSON.stringify({ const data = JSON.stringify({
id: null, id: null,
@@ -748,13 +758,13 @@ export class StratumV1Client {
} }
} }
private calculateDifficulty(header: Buffer): { submissionDifficulty: number, submissionHash: string } { private calculateDifficulty(header: Buffer): { submissionDifficulty: number, submissionHash: string, hashBuffer: Buffer } {
const hashResult = bitcoinjs.crypto.hash256(header); const hashResult = hash256(header);
const target = this.le256todouble(hashResult); const target = this.le256todouble(hashResult);
const submissionDifficulty = target === 0 ? Number.POSITIVE_INFINITY : TRUE_DIFF_ONE / target; const submissionDifficulty = target === 0 ? Number.POSITIVE_INFINITY : TRUE_DIFF_ONE / target;
return { submissionDifficulty, submissionHash: hashResult.toString('hex') }; return { submissionDifficulty, submissionHash: hashResult.toString('hex'), hashBuffer: hashResult };
} }
+98 -3
View File
@@ -1,4 +1,5 @@
import { ConfigService } from '@nestjs/config'; import { ConfigService } from '@nestjs/config';
import * as bitcoinjs from 'bitcoinjs-lib';
import { Socket } from 'net'; import { Socket } from 'net';
import { BehaviorSubject, firstValueFrom } from 'rxjs'; import { BehaviorSubject, firstValueFrom } from 'rxjs';
@@ -14,6 +15,7 @@ import {
serializeSubmitSharesExtended, serializeSubmitSharesExtended,
} from './sv2/sv2-extended-messages'; } from './sv2/sv2-extended-messages';
import { deserializeSetNewPrevHash, deserializeSubmitSharesError } from './sv2/sv2-messages'; import { deserializeSetNewPrevHash, deserializeSubmitSharesError } from './sv2/sv2-messages';
import { MiningJob } from './MiningJob';
import { StratumV2Client } from './StratumV2Client'; import { StratumV2Client } from './StratumV2Client';
describe('StratumV2Client extended channels', () => { describe('StratumV2Client extended channels', () => {
@@ -145,9 +147,87 @@ describe('StratumV2Client extended channels', () => {
expect(postAccountingPresenceUpdates.length).toBeGreaterThan(0); expect(postAccountingPresenceUpdates.length).toBeGreaterThan(0);
}); });
it('does not submit a block when only the reported SV2 difficulty is huge', async () => {
const { client, shareAccountingService, bitcoinRpcService, blocksService, jobTemplate } = await createClient();
(client as any).address = 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
(client as any).workerName = 'worker';
(client as any).sessionId = 'sv2-session';
(client as any).userAgent = 'test/sv2';
const job = new MiningJob(
bitcoinjs.networks.testnet,
'1',
[{ address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', percent: 100 }],
jobTemplate,
);
await (client as any).handleAcceptedShare(
{
nonce: 123,
ntime: parseInt(MockRecording1.TIME, 16),
version: jobTemplate.block.version,
},
{ extranoncePrefix: Buffer.from(MockRecording1.EXTRA_NONCE, 'hex') },
job,
jobTemplate,
Number.MAX_SAFE_INTEGER,
1024,
Buffer.alloc(32, 0xff),
);
expect(bitcoinRpcService.SUBMIT_BLOCK).not.toHaveBeenCalled();
expect(blocksService.save).not.toHaveBeenCalled();
expect(shareAccountingService.recordAcceptedShare).toHaveBeenCalledWith(expect.objectContaining({
submissionDifficulty: Number.MAX_SAFE_INTEGER,
isBlockCandidate: false,
}));
});
it('submits a block when the SV2 hash target exactly meets network target', async () => {
const { client, bitcoinRpcService, blocksService, notificationService, jobTemplate } = await createClient();
(client as any).address = 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
(client as any).workerName = 'worker';
(client as any).sessionId = 'sv2-session';
(client as any).userAgent = 'test/sv2';
const job = new MiningJob(
bitcoinjs.networks.testnet,
'1',
[{ address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', percent: 100 }],
jobTemplate,
);
await (client as any).handleAcceptedShare(
{
nonce: 123,
ntime: parseInt(MockRecording1.TIME, 16),
version: jobTemplate.block.version,
},
{ extranoncePrefix: Buffer.from(MockRecording1.EXTRA_NONCE, 'hex') },
job,
jobTemplate,
1,
1024,
Buffer.alloc(32),
);
expect(bitcoinRpcService.SUBMIT_BLOCK).toHaveBeenCalledWith(expect.any(String));
expect(blocksService.save).toHaveBeenCalledWith(expect.objectContaining({
height: MockRecording1.BLOCK_TEMPLATE.height,
minerAddress: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
worker: 'worker',
sessionId: 'sv2-session',
blockData: expect.any(String),
}));
expect(notificationService.notifySubscribersBlockFound).toHaveBeenCalled();
});
async function createClient(): Promise<{ async function createClient(): Promise<{
client: StratumV2Client; client: StratumV2Client;
sentFrames: any[]; sentFrames: any[];
bitcoinRpcService: { SUBMIT_BLOCK: jest.Mock };
blocksService: { save: jest.Mock };
notificationService: { notifySubscribersBlockFound: jest.Mock };
shareAccountingService: { recordAcceptedShare: jest.Mock }; shareAccountingService: { recordAcceptedShare: jest.Mock };
redisMessagingService: { setClientPresence: jest.Mock; removeClientPresence: jest.Mock }; redisMessagingService: { setClientPresence: jest.Mock; removeClientPresence: jest.Mock };
jobTemplate: any; jobTemplate: any;
@@ -196,6 +276,12 @@ describe('StratumV2Client extended channels', () => {
const shareAccountingService = { const shareAccountingService = {
recordAcceptedShare: jest.fn().mockResolvedValue(undefined), recordAcceptedShare: jest.fn().mockResolvedValue(undefined),
}; };
const notificationService = {
notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined),
};
const blocksService = {
save: jest.fn().mockResolvedValue(undefined),
};
const redisMessagingService = { const redisMessagingService = {
setClientPresence: jest.fn().mockResolvedValue(undefined), setClientPresence: jest.fn().mockResolvedValue(undefined),
removeClientPresence: jest.fn().mockResolvedValue(undefined), removeClientPresence: jest.fn().mockResolvedValue(undefined),
@@ -225,8 +311,8 @@ describe('StratumV2Client extended channels', () => {
stratumV1JobsService, stratumV1JobsService,
bitcoinRpcService as any, bitcoinRpcService as any,
clientService as any, clientService as any,
{ notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined) } as any, notificationService as any,
{ save: jest.fn().mockResolvedValue(undefined) } as any, blocksService as any,
{ {
get: jest.fn((key: string) => { get: jest.fn((key: string) => {
switch (key) { switch (key) {
@@ -251,6 +337,15 @@ describe('StratumV2Client extended channels', () => {
return Promise.resolve(); return Promise.resolve();
}); });
return { client, sentFrames, shareAccountingService, redisMessagingService, jobTemplate }; return {
client,
sentFrames,
bitcoinRpcService,
blocksService,
notificationService,
shareAccountingService,
redisMessagingService,
jobTemplate,
};
} }
}); });
+13 -5
View File
@@ -17,6 +17,7 @@ import { StratumV2Service } from '../services/stratum-v2.service';
import { IJobTemplate, StratumV1JobsService } from '../services/stratum-v1-jobs.service'; import { IJobTemplate, StratumV1JobsService } from '../services/stratum-v1-jobs.service';
import { patchCoinbasePrefixVarint } from '../utils/coinbase-prefix.utils'; import { patchCoinbasePrefixVarint } from '../utils/coinbase-prefix.utils';
import { DifficultyUtils } from '../utils/difficulty.utils'; import { DifficultyUtils } from '../utils/difficulty.utils';
import { hash256 } from '../utils/hash.utils';
import { MiningJob } from './MiningJob'; import { MiningJob } from './MiningJob';
import { StratumV1ClientStatistics } from './StratumV1ClientStatistics'; import { StratumV1ClientStatistics } from './StratumV1ClientStatistics';
import { TOTAL_EXTRANONCE_SIZE_BYTES } from './stratum.constants'; import { TOTAL_EXTRANONCE_SIZE_BYTES } from './stratum.constants';
@@ -560,7 +561,7 @@ export class StratumV2Client {
); );
channel.acceptedShareCount++; channel.acceptedShareCount++;
await this.handleAcceptedShare(submission, channel, job, jobTemplate, submissionDifficulty, jobDifficulty); await this.handleAcceptedShare(submission, channel, job, jobTemplate, submissionDifficulty, jobDifficulty, hashBuffer);
} }
private async handleSubmitSharesExtended(payload: Buffer): Promise<void> { private async handleSubmitSharesExtended(payload: Buffer): Promise<void> {
@@ -608,12 +609,12 @@ export class StratumV2Client {
submission.extranonce, submission.extranonce,
extendedJob.coinbaseSuffix, extendedJob.coinbaseSuffix,
]); ]);
let merkleRoot = bitcoinjs.crypto.hash256(coinbaseTxBytes); let merkleRoot = hash256(coinbaseTxBytes);
const merklePair = Buffer.alloc(64); const merklePair = Buffer.alloc(64);
for (const sibling of extendedJob.merklePath) { for (const sibling of extendedJob.merklePath) {
merklePair.set(merkleRoot, 0); merklePair.set(merkleRoot, 0);
merklePair.set(sibling, 32); merklePair.set(sibling, 32);
merkleRoot = bitcoinjs.crypto.hash256(merklePair); merkleRoot = hash256(merklePair);
} }
const header = this.buildHeader( const header = this.buildHeader(
@@ -650,7 +651,10 @@ export class StratumV2Client {
channel.acceptedShareCount++; channel.acceptedShareCount++;
let updatedJobBlock: bitcoinjs.Block = null; let updatedJobBlock: bitcoinjs.Block = null;
if (submissionDifficulty >= extendedJob.jobTemplate.blockData.networkDifficulty) { if (DifficultyUtils.meetsTarget(
hashBuffer,
DifficultyUtils.difficultyToTarget(extendedJob.jobTemplate.blockData.networkDifficulty),
)) {
updatedJobBlock = this.reconstructExtendedBlock(extendedJob, submission, merkleRoot, channel.extranoncePrefix); updatedJobBlock = this.reconstructExtendedBlock(extendedJob, submission, merkleRoot, channel.extranoncePrefix);
} }
await this.recordAcceptedShare(submissionDifficulty, jobDifficulty, extendedJob.jobTemplate, updatedJobBlock, { await this.recordAcceptedShare(submissionDifficulty, jobDifficulty, extendedJob.jobTemplate, updatedJobBlock, {
@@ -669,9 +673,13 @@ export class StratumV2Client {
jobTemplate: IJobTemplate, jobTemplate: IJobTemplate,
submissionDifficulty: number, submissionDifficulty: number,
jobDifficulty: number, jobDifficulty: number,
hashBuffer: Buffer,
): Promise<void> { ): Promise<void> {
let updatedJobBlock: bitcoinjs.Block = null; let updatedJobBlock: bitcoinjs.Block = null;
if (submissionDifficulty >= jobTemplate.blockData.networkDifficulty) { if (DifficultyUtils.meetsTarget(
hashBuffer,
DifficultyUtils.difficultyToTarget(jobTemplate.blockData.networkDifficulty),
)) {
const versionMask = submission.version ^ jobTemplate.block.version; const versionMask = submission.version ^ jobTemplate.block.version;
updatedJobBlock = job.copyAndUpdateBlock( updatedJobBlock = job.copyAndUpdateBlock(
jobTemplate, jobTemplate,
@@ -37,6 +37,10 @@ describe('MiningSubmitMessage', () => {
expect(errors).toEqual([]); expect(errors).toEqual([]);
}); });
it('should hash submissions deterministically', () => {
expect(message.hash()).toBe('t2bFzhZ6mketRxa5nOoKrNxG3RkFIZuOwY1WewFtv9k=');
});
it('should reject short extranonce2 submissions', async () => { it('should reject short extranonce2 submissions', async () => {
const shortMessage = plainToInstance( const shortMessage = plainToInstance(
MiningSubmitMessage, MiningSubmitMessage,
@@ -4,7 +4,7 @@ import { ArrayMaxSize, ArrayMinSize, IsArray, IsString, Length } from 'class-val
import { eRequestMethod } from '../enums/eRequestMethod'; import { eRequestMethod } from '../enums/eRequestMethod';
import { EXTRANONCE2_SIZE_BYTES } from '../stratum.constants'; import { EXTRANONCE2_SIZE_BYTES } from '../stratum.constants';
import { StratumBaseMessage } from './StratumBaseMessage'; import { StratumBaseMessage } from './StratumBaseMessage';
import * as bitcoinjs from 'bitcoinjs-lib'; import { hash256 } from '../../utils/hash.utils';
export class MiningSubmitMessage extends StratumBaseMessage { export class MiningSubmitMessage extends StratumBaseMessage {
@@ -69,7 +69,7 @@ export class MiningSubmitMessage extends StratumBaseMessage {
public hash(): string{ public hash(): string{
const buffer = Buffer.from(this.versionMask + this.nonce + this.extraNonce2 + this.ntime + this.jobId); const buffer = Buffer.from(this.versionMask + this.nonce + this.extraNonce2 + this.ntime + this.jobId);
return bitcoinjs.crypto.hash256(buffer).toString('base64'); return hash256(buffer).toString('base64');
} }
+2 -1
View File
@@ -5,6 +5,7 @@ import * as merkleProof from 'merkle-lib/proof';
import { combineLatest, delay, filter, from, interval, map, Observable, shareReplay, startWith, switchMap, tap } from 'rxjs'; import { combineLatest, delay, filter, from, interval, map, Observable, shareReplay, startWith, switchMap, tap } from 'rxjs';
import { MiningJob } from '../models/MiningJob'; import { MiningJob } from '../models/MiningJob';
import { hash256 } from '../utils/hash.utils';
import { BitcoinRpcService } from './bitcoin-rpc.service'; import { BitcoinRpcService } from './bitcoin-rpc.service';
export interface IJobTemplate { export interface IJobTemplate {
@@ -95,7 +96,7 @@ export class StratumV1JobsService {
const transactionBuffers = transactions.map(tx => tx.getHash(false)); const transactionBuffers = transactions.map(tx => tx.getHash(false));
const merkleTree = merkle(transactionBuffers, bitcoinjs.crypto.hash256); const merkleTree = merkle(transactionBuffers, hash256);
const merkleBranches: Buffer[] = merkleProof(merkleTree, transactionBuffers[0]).filter(h => h != null); const merkleBranches: Buffer[] = merkleProof(merkleTree, transactionBuffers[0]).filter(h => h != null);
block.merkleRoot = merkleBranches.pop(); block.merkleRoot = merkleBranches.pop();
+62
View File
@@ -0,0 +1,62 @@
import { DifficultyUtils } from './difficulty.utils';
describe('DifficultyUtils', () => {
function incrementLe256(target: Buffer): Buffer {
const next = Buffer.from(target);
for (let i = 0; i < next.length; i++) {
next[i]++;
if (next[i] !== 0) {
break;
}
}
return next;
}
it('round-trips common difficulty values through compact targets', () => {
for (const difficulty of [0.001, 1, 128, 512, 1000, 1_000_000, 1_000_000_000, 1e30]) {
const target = DifficultyUtils.difficultyToTarget(difficulty);
const roundTrip = DifficultyUtils.targetToDifficulty(target);
expect(Math.abs(roundTrip - difficulty) / difficulty).toBeLessThan(1e-12);
}
});
it('compares little-endian targets without converting to BigInt', () => {
const target = Buffer.alloc(32);
target[0] = 0x80;
target[30] = 0x01;
const easier = Buffer.from(target);
const harder = Buffer.from(target);
easier[31] += 1;
harder[0] -= 1;
expect(DifficultyUtils.meetsTarget(target, target)).toBe(true);
expect(DifficultyUtils.meetsTarget(harder, target)).toBe(true);
expect(DifficultyUtils.meetsTarget(easier, target)).toBe(false);
});
it('calculates difficulty consistently with the returned hash buffer', () => {
const header = Buffer.alloc(80, 1);
const result = DifficultyUtils.calculateDifficulty(header);
expect(result.submissionDifficulty).toBeCloseTo(
DifficultyUtils.targetToDifficulty(result.hashBuffer),
10,
);
expect(result.submissionHash).toBe(result.hashBuffer.toString('hex'));
});
it('keeps huge-difficulty boundary checks exact with target bytes', () => {
const target = DifficultyUtils.difficultyToTarget(1e30);
const justTooEasy = incrementLe256(target);
expect(DifficultyUtils.targetToDifficulty(target)).toBeGreaterThan(1e29);
expect(DifficultyUtils.meetsTarget(target, target)).toBe(true);
expect(DifficultyUtils.meetsTarget(justTooEasy, target)).toBe(false);
});
it('treats a zero hash target as infinite reported difficulty', () => {
expect(DifficultyUtils.targetToDifficulty(Buffer.alloc(32))).toBe(Number.POSITIVE_INFINITY);
});
});
+35 -26
View File
@@ -1,35 +1,25 @@
import * as bitcoinjs from 'bitcoinjs-lib'; import { hash256 } from './hash.utils';
const TRUE_DIFF_ONE_BIGINT = BigInt( const TRUE_DIFF_ONE_BIGINT = BigInt(
'26959535291011309493156476344723991336010898738574164086137773096960', '26959535291011309493156476344723991336010898738574164086137773096960',
); );
const TRUE_DIFF_ONE_NUMBER = 2.695953529101131e67;
const TWO_TO_256 = 1n << 256n; const TWO_TO_256 = 1n << 256n;
const FRACTION_SCALE = 1_000_000_000_000_000n;
const FRACTION_SCALE_NUM = 1e15;
function bigIntRatioToDifficulty(divisor: bigint): number {
if (divisor === 0n) {
return Number.POSITIVE_INFINITY;
}
const scaled = (TRUE_DIFF_ONE_BIGINT * FRACTION_SCALE) / divisor;
return Number(scaled) / FRACTION_SCALE_NUM;
}
export class DifficultyUtils { export class DifficultyUtils {
public static calculateDifficulty(header: Buffer): { submissionDifficulty: number; submissionHash: string; hashBuffer: Buffer } { public static calculateDifficulty(header: Buffer): { submissionDifficulty: number; submissionHash: string; hashBuffer: Buffer } {
const hashResult = bitcoinjs.crypto.hash256(header); const hashResult = hash256(header);
const target = DifficultyUtils.le256ToBigInt(hashResult); const target = DifficultyUtils.le256ToDouble(hashResult);
return { return {
submissionDifficulty: bigIntRatioToDifficulty(target), submissionDifficulty: target === 0 ? Number.POSITIVE_INFINITY : TRUE_DIFF_ONE_NUMBER / target,
submissionHash: hashResult.toString('hex'), submissionHash: hashResult.toString('hex'),
hashBuffer: hashResult, hashBuffer: hashResult,
}; };
} }
public static meetsTarget(hashBuffer: Buffer, target: Buffer): boolean { public static meetsTarget(hashBuffer: Buffer, target: Buffer): boolean {
return DifficultyUtils.le256ToBigInt(hashBuffer) <= DifficultyUtils.le256ToBigInt(target); return DifficultyUtils.compareLe256(hashBuffer, target) <= 0;
} }
public static difficultyToTarget(difficulty: number): Buffer { public static difficultyToTarget(difficulty: number): Buffer {
@@ -51,8 +41,8 @@ export class DifficultyUtils {
throw new Error('Target must be 32 bytes'); throw new Error('Target must be 32 bytes');
} }
const targetBigInt = DifficultyUtils.le256ToBigInt(target); const targetNumber = DifficultyUtils.le256ToDouble(target);
return bigIntRatioToDifficulty(targetBigInt); return targetNumber === 0 ? Number.POSITIVE_INFINITY : TRUE_DIFF_ONE_NUMBER / targetNumber;
} }
public static hashRateToDifficulty(hashRate: number, sharesPerMinute: number): number { public static hashRateToDifficulty(hashRate: number, sharesPerMinute: number): number {
@@ -65,15 +55,11 @@ export class DifficultyUtils {
return difficulty; return difficulty;
} }
const maxTargetBigInt = DifficultyUtils.le256ToBigInt(maxTarget); if (DifficultyUtils.isZeroTarget(maxTarget)) {
if (maxTargetBigInt === 0n) {
return difficulty; return difficulty;
} }
const computedTargetBigInt = DifficultyUtils.le256ToBigInt( if (DifficultyUtils.compareLe256(DifficultyUtils.difficultyToTarget(difficulty), maxTarget) > 0) {
DifficultyUtils.difficultyToTarget(difficulty),
);
if (computedTargetBigInt > maxTargetBigInt) {
const clamped = DifficultyUtils.targetToDifficulty(maxTarget); const clamped = DifficultyUtils.targetToDifficulty(maxTarget);
return Number.isFinite(clamped) && clamped > 0 ? clamped : difficulty; return Number.isFinite(clamped) && clamped > 0 ? clamped : difficulty;
} }
@@ -111,7 +97,30 @@ export class DifficultyUtils {
return buf; return buf;
} }
private static le256ToBigInt(target: Buffer): bigint { private static le256ToDouble(target: Buffer): number {
return target.reduceRight((acc, byte) => (acc << 8n) | BigInt(byte), 0n); let value = 0;
for (let i = target.length - 1; i >= 0; i--) {
value = value * 256 + target[i];
}
return value;
}
private static compareLe256(left: Buffer, right: Buffer): number {
for (let i = 31; i >= 0; i--) {
const diff = left[i] - right[i];
if (diff !== 0) {
return diff;
}
}
return 0;
}
private static isZeroTarget(target: Buffer): boolean {
for (let i = 0; i < target.length; i++) {
if (target[i] !== 0) {
return false;
}
}
return true;
} }
} }
+15
View File
@@ -0,0 +1,15 @@
import * as crypto from 'crypto';
import { hash256 } from './hash.utils';
describe('hash256', () => {
it('matches native double sha256 for known input', () => {
const data = Buffer.from('public-pool-share-validation', 'utf8');
const expected = crypto
.createHash('sha256')
.update(crypto.createHash('sha256').update(data).digest())
.digest('hex');
expect(hash256(data).toString('hex')).toBe(expected);
});
});
+6
View File
@@ -0,0 +1,6 @@
import * as crypto from 'crypto';
export function hash256(data: Buffer): Buffer {
const first = crypto.createHash('sha256').update(data).digest();
return crypto.createHash('sha256').update(first).digest();
}
+107
View File
@@ -42,6 +42,8 @@ describe('TimescaleDB and Redis integration', () => {
}); });
beforeEach(async () => { 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 accepted_share_entity`);
await dataSource.query(`DELETE FROM client_entity`); await dataSource.query(`DELETE FROM client_entity`);
await dataSource.query(`REFRESH MATERIALIZED VIEW user_agent_report_view`); 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', sequence_name: 'accepted_share_index_seq',
index_name: '"IDX_accepted_share_order"', 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 () => { 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 () => { it('should exclude soft-deleted clients from the user-agent report', async () => {
const userAgent = `integration-reconnect-${Date.now()}`; const userAgent = `integration-reconnect-${Date.now()}`;
const repository = dataSource.getRepository(ClientEntity); const repository = dataSource.getRepository(ClientEntity);