mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 09:05:06 -07:00
Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bcd6bb0aa7 | ||
|
|
1c5a3d9a94 |
@@ -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
|
||||
|
||||
@@ -20,6 +20,7 @@ lerna-debug.log*
|
||||
/.nyc_output
|
||||
/.baseline-*.json
|
||||
/.test-artifacts
|
||||
/.profiles
|
||||
|
||||
# IDEs and editors
|
||||
/.idea
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -2,6 +2,7 @@ import { AddressType, getAddressInfo } from 'bitcoin-address-validation';
|
||||
import * as bitcoinjs from 'bitcoinjs-lib';
|
||||
|
||||
import { IJobTemplate } from '../services/stratum-v1-jobs.service';
|
||||
import { hash256 } from '../utils/hash.utils';
|
||||
import { eResponseMethod } from './enums/eResponseMethod';
|
||||
import { IMiningNotify } from './stratum-messages/IMiningNotify';
|
||||
import { TOTAL_EXTRANONCE_SIZE_BYTES } from './stratum.constants';
|
||||
@@ -96,7 +97,7 @@ export class MiningJob {
|
||||
Buffer.from(`${extraNonce}${extraNonce2}`, 'hex'),
|
||||
this.coinbasePart2Buffer,
|
||||
]);
|
||||
const coinbaseHash = bitcoinjs.crypto.hash256(coinbaseBuffer);
|
||||
const coinbaseHash = hash256(coinbaseBuffer);
|
||||
const merkleRoot = this.calculateMerkleRootHash(coinbaseHash, this.merkleBranchBuffers);
|
||||
|
||||
let version = jobTemplate.block.version;
|
||||
@@ -154,7 +155,7 @@ export class MiningJob {
|
||||
|
||||
for (let i = 0; i < merkleBranches.length; i++) {
|
||||
bothMerkles.set(merkleBranches[i], 32);
|
||||
newRoot = bitcoinjs.crypto.hash256(bothMerkles);
|
||||
newRoot = hash256(bothMerkles);
|
||||
bothMerkles.set(newRoot);
|
||||
}
|
||||
|
||||
|
||||
@@ -9,6 +9,7 @@ import { ClientService } from '../ORM/client/client.service';
|
||||
import { BitcoinRpcService as MockBitcoinRpcService } from '../services/bitcoin-rpc.service';
|
||||
import { NotificationService } from '../services/notification.service';
|
||||
import { StratumV1JobsService } from '../services/stratum-v1-jobs.service';
|
||||
import { DifficultyUtils } from '../utils/difficulty.utils';
|
||||
import { IBlockTemplate } from './bitcoin-rpc/IBlockTemplate';
|
||||
import { MiningJob } from './MiningJob';
|
||||
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, 'calculateDifficulty').mockReturnValue({
|
||||
submissionDifficulty: 1024,
|
||||
submissionHash: 'share'
|
||||
submissionHash: 'share',
|
||||
hashBuffer: DifficultyUtils.difficultyToTarget(1024),
|
||||
});
|
||||
const getSettingsSpy = jest.spyOn(addressSettings, 'getSettings');
|
||||
const updateIfHigherSpy = jest.spyOn(addressSettings as any, 'updateBestDifficultyIfHigher').mockResolvedValue({ affected: 1 });
|
||||
@@ -495,6 +497,27 @@ describe('StratumV1Client', () => {
|
||||
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 () => {
|
||||
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, 'calculateDifficulty').mockReturnValue({
|
||||
submissionDifficulty: Number.MAX_SAFE_INTEGER,
|
||||
submissionHash: 'block-share'
|
||||
submissionHash: 'block-share',
|
||||
hashBuffer: Buffer.alloc(32),
|
||||
});
|
||||
jest.spyOn(addressSettings, 'resetBestDifficultyAndShares').mockResolvedValue(undefined);
|
||||
|
||||
|
||||
@@ -16,6 +16,8 @@ import { BitcoinRpcService } from '../services/bitcoin-rpc.service';
|
||||
import { NotificationService } from '../services/notification.service';
|
||||
import { RedisMessagingService } from '../services/redis-messaging.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 { eResponseMethod } from './enums/eResponseMethod';
|
||||
import { eStratumErrorCode } from './enums/eStratumErrorCode';
|
||||
@@ -51,6 +53,7 @@ export class StratumV1Client {
|
||||
private stratumInitialized = false;
|
||||
private usedSuggestedDifficulty = false;
|
||||
private sessionDifficulty: number = 100000;
|
||||
private sessionDifficultyTarget: Buffer = DifficultyUtils.difficultyToTarget(this.sessionDifficulty);
|
||||
|
||||
private clientEntity: ClientEntity;
|
||||
private creatingEntity: Promise<void>;
|
||||
@@ -268,6 +271,7 @@ export class StratumV1Client {
|
||||
this.clientAuthorization = authorizationMessage;
|
||||
if (this.clientSuggestedDifficulty == null && this.clientAuthorization.startingDiff != null && this.clientAuthorization.startingDiff > this.sessionDifficulty) {
|
||||
this.sessionDifficulty = this.clientAuthorization.startingDiff;
|
||||
this.sessionDifficultyTarget = DifficultyUtils.difficultyToTarget(this.sessionDifficulty);
|
||||
}
|
||||
const success = await this.write(JSON.stringify(this.clientAuthorization.response()) + '\n');
|
||||
if (!success) {
|
||||
@@ -309,6 +313,7 @@ export class StratumV1Client {
|
||||
|
||||
this.clientSuggestedDifficulty = suggestDifficultyMessage;
|
||||
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');
|
||||
if (!success) {
|
||||
return;
|
||||
@@ -610,19 +615,23 @@ export class StratumV1Client {
|
||||
submission.extraNonce2,
|
||||
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}`);
|
||||
|
||||
|
||||
if (submissionDifficulty >= this.sessionDifficulty) {
|
||||
if (DifficultyUtils.meetsTarget(hashBuffer, this.sessionDifficultyTarget)) {
|
||||
const success = await this.write(JSON.stringify(submission.response()) + '\n');
|
||||
if (!success) {
|
||||
return false;
|
||||
}
|
||||
|
||||
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 !!!');
|
||||
const updatedJobBlock = job.copyAndUpdateBlock(
|
||||
jobTemplate,
|
||||
@@ -668,7 +677,7 @@ export class StratumV1Client {
|
||||
? (jobTemplate.block.version ^ versionMask).toString(16)
|
||||
: jobTemplate.block.version.toString(16),
|
||||
extraNonce2: submission.extraNonce2,
|
||||
isBlockCandidate: submissionDifficulty >= jobTemplate.blockData.networkDifficulty,
|
||||
isBlockCandidate,
|
||||
blockSubmissionResult,
|
||||
});
|
||||
await this.statistics.addShares(this.clientEntity, this.sessionDifficulty);
|
||||
@@ -718,6 +727,7 @@ export class StratumV1Client {
|
||||
if (targetDiff != this.sessionDifficulty) {
|
||||
//console.log(`Adjusting ${this.extraNonceAndSessionId} difficulty from ${this.sessionDifficulty} to ${targetDiff}`);
|
||||
this.sessionDifficulty = targetDiff;
|
||||
this.sessionDifficultyTarget = DifficultyUtils.difficultyToTarget(this.sessionDifficulty);
|
||||
|
||||
const data = JSON.stringify({
|
||||
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 submissionDifficulty = target === 0 ? Number.POSITIVE_INFINITY : TRUE_DIFF_ONE / target;
|
||||
return { submissionDifficulty, submissionHash: hashResult.toString('hex') };
|
||||
return { submissionDifficulty, submissionHash: hashResult.toString('hex'), hashBuffer: hashResult };
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import { ConfigService } from '@nestjs/config';
|
||||
import * as bitcoinjs from 'bitcoinjs-lib';
|
||||
import { Socket } from 'net';
|
||||
import { BehaviorSubject, firstValueFrom } from 'rxjs';
|
||||
|
||||
@@ -14,6 +15,7 @@ import {
|
||||
serializeSubmitSharesExtended,
|
||||
} from './sv2/sv2-extended-messages';
|
||||
import { deserializeSetNewPrevHash, deserializeSubmitSharesError } from './sv2/sv2-messages';
|
||||
import { MiningJob } from './MiningJob';
|
||||
import { StratumV2Client } from './StratumV2Client';
|
||||
|
||||
describe('StratumV2Client extended channels', () => {
|
||||
@@ -145,9 +147,87 @@ describe('StratumV2Client extended channels', () => {
|
||||
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<{
|
||||
client: StratumV2Client;
|
||||
sentFrames: any[];
|
||||
bitcoinRpcService: { SUBMIT_BLOCK: jest.Mock };
|
||||
blocksService: { save: jest.Mock };
|
||||
notificationService: { notifySubscribersBlockFound: jest.Mock };
|
||||
shareAccountingService: { recordAcceptedShare: jest.Mock };
|
||||
redisMessagingService: { setClientPresence: jest.Mock; removeClientPresence: jest.Mock };
|
||||
jobTemplate: any;
|
||||
@@ -196,6 +276,12 @@ describe('StratumV2Client extended channels', () => {
|
||||
const shareAccountingService = {
|
||||
recordAcceptedShare: jest.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
const notificationService = {
|
||||
notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
const blocksService = {
|
||||
save: jest.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
const redisMessagingService = {
|
||||
setClientPresence: jest.fn().mockResolvedValue(undefined),
|
||||
removeClientPresence: jest.fn().mockResolvedValue(undefined),
|
||||
@@ -225,8 +311,8 @@ describe('StratumV2Client extended channels', () => {
|
||||
stratumV1JobsService,
|
||||
bitcoinRpcService as any,
|
||||
clientService as any,
|
||||
{ notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined) } as any,
|
||||
{ save: jest.fn().mockResolvedValue(undefined) } as any,
|
||||
notificationService as any,
|
||||
blocksService as any,
|
||||
{
|
||||
get: jest.fn((key: string) => {
|
||||
switch (key) {
|
||||
@@ -251,6 +337,15 @@ describe('StratumV2Client extended channels', () => {
|
||||
return Promise.resolve();
|
||||
});
|
||||
|
||||
return { client, sentFrames, shareAccountingService, redisMessagingService, jobTemplate };
|
||||
return {
|
||||
client,
|
||||
sentFrames,
|
||||
bitcoinRpcService,
|
||||
blocksService,
|
||||
notificationService,
|
||||
shareAccountingService,
|
||||
redisMessagingService,
|
||||
jobTemplate,
|
||||
};
|
||||
}
|
||||
});
|
||||
|
||||
@@ -17,6 +17,7 @@ import { StratumV2Service } from '../services/stratum-v2.service';
|
||||
import { IJobTemplate, StratumV1JobsService } from '../services/stratum-v1-jobs.service';
|
||||
import { patchCoinbasePrefixVarint } from '../utils/coinbase-prefix.utils';
|
||||
import { DifficultyUtils } from '../utils/difficulty.utils';
|
||||
import { hash256 } from '../utils/hash.utils';
|
||||
import { MiningJob } from './MiningJob';
|
||||
import { StratumV1ClientStatistics } from './StratumV1ClientStatistics';
|
||||
import { TOTAL_EXTRANONCE_SIZE_BYTES } from './stratum.constants';
|
||||
@@ -560,7 +561,7 @@ export class StratumV2Client {
|
||||
);
|
||||
|
||||
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> {
|
||||
@@ -608,12 +609,12 @@ export class StratumV2Client {
|
||||
submission.extranonce,
|
||||
extendedJob.coinbaseSuffix,
|
||||
]);
|
||||
let merkleRoot = bitcoinjs.crypto.hash256(coinbaseTxBytes);
|
||||
let merkleRoot = hash256(coinbaseTxBytes);
|
||||
const merklePair = Buffer.alloc(64);
|
||||
for (const sibling of extendedJob.merklePath) {
|
||||
merklePair.set(merkleRoot, 0);
|
||||
merklePair.set(sibling, 32);
|
||||
merkleRoot = bitcoinjs.crypto.hash256(merklePair);
|
||||
merkleRoot = hash256(merklePair);
|
||||
}
|
||||
|
||||
const header = this.buildHeader(
|
||||
@@ -650,7 +651,10 @@ export class StratumV2Client {
|
||||
|
||||
channel.acceptedShareCount++;
|
||||
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);
|
||||
}
|
||||
await this.recordAcceptedShare(submissionDifficulty, jobDifficulty, extendedJob.jobTemplate, updatedJobBlock, {
|
||||
@@ -669,9 +673,13 @@ export class StratumV2Client {
|
||||
jobTemplate: IJobTemplate,
|
||||
submissionDifficulty: number,
|
||||
jobDifficulty: number,
|
||||
hashBuffer: Buffer,
|
||||
): Promise<void> {
|
||||
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;
|
||||
updatedJobBlock = job.copyAndUpdateBlock(
|
||||
jobTemplate,
|
||||
|
||||
@@ -37,6 +37,10 @@ describe('MiningSubmitMessage', () => {
|
||||
expect(errors).toEqual([]);
|
||||
});
|
||||
|
||||
it('should hash submissions deterministically', () => {
|
||||
expect(message.hash()).toBe('t2bFzhZ6mketRxa5nOoKrNxG3RkFIZuOwY1WewFtv9k=');
|
||||
});
|
||||
|
||||
it('should reject short extranonce2 submissions', async () => {
|
||||
const shortMessage = plainToInstance(
|
||||
MiningSubmitMessage,
|
||||
|
||||
@@ -4,7 +4,7 @@ import { ArrayMaxSize, ArrayMinSize, IsArray, IsString, Length } from 'class-val
|
||||
import { eRequestMethod } from '../enums/eRequestMethod';
|
||||
import { EXTRANONCE2_SIZE_BYTES } from '../stratum.constants';
|
||||
import { StratumBaseMessage } from './StratumBaseMessage';
|
||||
import * as bitcoinjs from 'bitcoinjs-lib';
|
||||
import { hash256 } from '../../utils/hash.utils';
|
||||
|
||||
|
||||
export class MiningSubmitMessage extends StratumBaseMessage {
|
||||
@@ -69,7 +69,7 @@ export class MiningSubmitMessage extends StratumBaseMessage {
|
||||
|
||||
public hash(): string{
|
||||
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');
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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 { MiningJob } from '../models/MiningJob';
|
||||
import { hash256 } from '../utils/hash.utils';
|
||||
import { BitcoinRpcService } from './bitcoin-rpc.service';
|
||||
|
||||
export interface IJobTemplate {
|
||||
@@ -95,7 +96,7 @@ export class StratumV1JobsService {
|
||||
|
||||
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);
|
||||
block.merkleRoot = merkleBranches.pop();
|
||||
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -1,35 +1,25 @@
|
||||
import * as bitcoinjs from 'bitcoinjs-lib';
|
||||
import { hash256 } from './hash.utils';
|
||||
|
||||
const TRUE_DIFF_ONE_BIGINT = BigInt(
|
||||
'26959535291011309493156476344723991336010898738574164086137773096960',
|
||||
);
|
||||
const TRUE_DIFF_ONE_NUMBER = 2.695953529101131e67;
|
||||
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 {
|
||||
public static calculateDifficulty(header: Buffer): { submissionDifficulty: number; submissionHash: string; hashBuffer: Buffer } {
|
||||
const hashResult = bitcoinjs.crypto.hash256(header);
|
||||
const target = DifficultyUtils.le256ToBigInt(hashResult);
|
||||
const hashResult = hash256(header);
|
||||
const target = DifficultyUtils.le256ToDouble(hashResult);
|
||||
|
||||
return {
|
||||
submissionDifficulty: bigIntRatioToDifficulty(target),
|
||||
submissionDifficulty: target === 0 ? Number.POSITIVE_INFINITY : TRUE_DIFF_ONE_NUMBER / target,
|
||||
submissionHash: hashResult.toString('hex'),
|
||||
hashBuffer: hashResult,
|
||||
};
|
||||
}
|
||||
|
||||
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 {
|
||||
@@ -51,8 +41,8 @@ export class DifficultyUtils {
|
||||
throw new Error('Target must be 32 bytes');
|
||||
}
|
||||
|
||||
const targetBigInt = DifficultyUtils.le256ToBigInt(target);
|
||||
return bigIntRatioToDifficulty(targetBigInt);
|
||||
const targetNumber = DifficultyUtils.le256ToDouble(target);
|
||||
return targetNumber === 0 ? Number.POSITIVE_INFINITY : TRUE_DIFF_ONE_NUMBER / targetNumber;
|
||||
}
|
||||
|
||||
public static hashRateToDifficulty(hashRate: number, sharesPerMinute: number): number {
|
||||
@@ -65,15 +55,11 @@ export class DifficultyUtils {
|
||||
return difficulty;
|
||||
}
|
||||
|
||||
const maxTargetBigInt = DifficultyUtils.le256ToBigInt(maxTarget);
|
||||
if (maxTargetBigInt === 0n) {
|
||||
if (DifficultyUtils.isZeroTarget(maxTarget)) {
|
||||
return difficulty;
|
||||
}
|
||||
|
||||
const computedTargetBigInt = DifficultyUtils.le256ToBigInt(
|
||||
DifficultyUtils.difficultyToTarget(difficulty),
|
||||
);
|
||||
if (computedTargetBigInt > maxTargetBigInt) {
|
||||
if (DifficultyUtils.compareLe256(DifficultyUtils.difficultyToTarget(difficulty), maxTarget) > 0) {
|
||||
const clamped = DifficultyUtils.targetToDifficulty(maxTarget);
|
||||
return Number.isFinite(clamped) && clamped > 0 ? clamped : difficulty;
|
||||
}
|
||||
@@ -111,7 +97,30 @@ export class DifficultyUtils {
|
||||
return buf;
|
||||
}
|
||||
|
||||
private static le256ToBigInt(target: Buffer): bigint {
|
||||
return target.reduceRight((acc, byte) => (acc << 8n) | BigInt(byte), 0n);
|
||||
private static le256ToDouble(target: Buffer): number {
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -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();
|
||||
}
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user