From d9fa933825202c1158a485cb2b10e394522652ad Mon Sep 17 00:00:00 2001 From: Ben Date: Wed, 6 May 2026 20:58:08 -0400 Subject: [PATCH] optimize stratum broadcast fanout --- src/models/MiningJob.spec.ts | 26 +++ src/models/MiningJob.ts | 17 ++ src/models/StratumV1Client.spec.ts | 166 +++++++++++++++- src/models/StratumV1Client.ts | 298 ++++++++++++++++++++++++----- 4 files changed, 452 insertions(+), 55 deletions(-) diff --git a/src/models/MiningJob.spec.ts b/src/models/MiningJob.spec.ts index 4c7e2a9..4330b7c 100644 --- a/src/models/MiningJob.spec.ts +++ b/src/models/MiningJob.spec.ts @@ -123,6 +123,32 @@ describe('MiningJob', () => { expect(clone.ins[0].witness[0]).toHaveLength(32); }); + it('should clone cached job data with a new job id', () => { + const clone = job.cloneForJobId('2'); + const originalNotify = JSON.parse(job.response(jobTemplate)); + const cloneNotify = JSON.parse(clone.response(jobTemplate)); + + expect(originalNotify.params[0]).toBe('1'); + expect(cloneNotify.params[0]).toBe('2'); + expect(cloneNotify.params[2]).toBe(originalNotify.params[2]); + expect(cloneNotify.params[3]).toBe(originalNotify.params[3]); + expect(clone.buildHeaderBuffer( + jobTemplate, + parseInt('00002000', 16), + parseInt('ed460d91', 16), + '57a6f098', + 'c708000000000000', + parseInt(MockRecording1.TIME, 16) + ).equals(job.buildHeaderBuffer( + jobTemplate, + parseInt('00002000', 16), + parseInt('ed460d91', 16), + '57a6f098', + 'c708000000000000', + parseInt(MockRecording1.TIME, 16) + ))).toBe(true); + }); + it('should omit oversized pool identifiers from the coinbase script', () => { jest.spyOn(console, 'warn').mockImplementation(() => undefined); const oversizedIdentifier = 'x'.repeat(120); diff --git a/src/models/MiningJob.ts b/src/models/MiningJob.ts index 6c73832..3d67816 100644 --- a/src/models/MiningJob.ts +++ b/src/models/MiningJob.ts @@ -98,6 +98,23 @@ export class MiningJob { return this.coinbaseTransaction.__toBuffer().toString('hex'); } + public cloneForJobId(jobId: string): MiningJob { + const job = Object.create(MiningJob.prototype) as MiningJob; + job.network = this.network; + job.jobId = jobId; + job.coinbaseTransaction = this.coinbaseTransaction; + job.coinbasePart1 = this.coinbasePart1; + job.coinbasePart2 = this.coinbasePart2; + job.coinbasePart1Buffer = this.coinbasePart1Buffer; + job.coinbasePart2Buffer = this.coinbasePart2Buffer; + job.merkleBranchBuffers = this.merkleBranchBuffers; + job.jobTemplateId = this.jobTemplateId; + job.networkDifficulty = this.networkDifficulty; + job.creation = new Date().getTime(); + job.retiredAt = undefined; + return job; + } + public getCoinbasePrefixBuffer(): Buffer { return Buffer.from(this.coinbasePart1Buffer); } diff --git a/src/models/StratumV1Client.spec.ts b/src/models/StratumV1Client.spec.ts index 4d9ef98..3629035 100644 --- a/src/models/StratumV1Client.spec.ts +++ b/src/models/StratumV1Client.spec.ts @@ -104,6 +104,8 @@ describe('StratumV1Client', () => { jest.useFakeTimers({ advanceTimers: true }) jest.setSystemTime(new Date(parseInt(MockRecording1.TIME, 16) * 1000)); + delete process.env.STRATUM_DIFFICULTY_CHECK_INTERVAL_MS; + delete process.env.STRATUM_DIFFICULTY_CHECK_BATCH_SIZE; consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined); consoleErrorSpy = jest.spyOn(console, 'error').mockImplementation(() => undefined); consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined); @@ -131,9 +133,24 @@ describe('StratumV1Client', () => { }); (StratumV1Client as any).blockedUserAgentLogState.clear(); (StratumV1Client as any).validationErrorLogState.clear(); - (StratumV1Client as any).jobBroadcastQueue = []; + (StratumV1Client as any).activeClients.clear(); + (StratumV1Client as any).jobBroadcastSubscription?.unsubscribe(); + (StratumV1Client as any).jobBroadcastSubscription = null; + (StratumV1Client as any).latestJobTemplate = null; + (StratumV1Client as any).pendingJobTemplate = null; + (StratumV1Client as any).jobBroadcastIterator = null; (StratumV1Client as any).jobBroadcastDraining = false; - (StratumV1Client as any).jobBroadcastDrainOffset = 0; + (StratumV1Client as any).clientJobQueue = []; + (StratumV1Client as any).clientJobDraining = false; + (StratumV1Client as any).clientJobDrainOffset = 0; + if ((StratumV1Client as any).difficultyCheckInterval != null) { + clearInterval((StratumV1Client as any).difficultyCheckInterval); + } + (StratumV1Client as any).difficultyCheckInterval = null; + (StratumV1Client as any).difficultyCheckIterator = null; + (StratumV1Client as any).difficultyCheckDraining = false; + (StratumV1Client as any).miningJobCacheTemplateId = null; + (StratumV1Client as any).miningJobCache.clear(); bitcoinRpcService = { newBlockTemplate$: newBlockEmitter.asObservable(), @@ -432,11 +449,10 @@ describe('StratumV1Client', () => { }); - it('should queue mining job broadcasts instead of building jobs in the observable callback', async () => { + it('should broadcast mining jobs through the shared active-client broadcaster', async () => { jest.useFakeTimers(); jest.setSystemTime(new Date(parseInt(MockRecording1.TIME, 16) * 1000)); jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); - const enqueueNewMiningJobSpy = jest.spyOn(client as any, 'enqueueNewMiningJob'); const sendNewMiningJobSpy = jest.spyOn(client as any, 'sendNewMiningJob'); (client as any).clientSubscription = { userAgent: 'bitaxe' }; @@ -449,7 +465,8 @@ describe('StratumV1Client', () => { await Promise.resolve(); await Promise.resolve(); - expect(enqueueNewMiningJobSpy).toHaveBeenCalled(); + expect((StratumV1Client as any).activeClients.has(client)).toBe(true); + expect((StratumV1Client as any).jobBroadcastSubscription).not.toBeNull(); jest.advanceTimersByTime(0); await Promise.resolve(); @@ -458,6 +475,145 @@ describe('StratumV1Client', () => { expect(sendNewMiningJobSpy).toHaveBeenCalled(); }); + it('should send the current job only to newly initialized clients after the shared broadcaster is running', async () => { + jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); + const firstClientSendSpy = jest.spyOn(client as any, 'sendNewMiningJob'); + + (client as any).clientSubscription = { userAgent: 'bitaxe' }; + (client as any).clientAuthorization = { + address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', + worker: 'bitaxe3' + }; + (client as any).statistics = {}; + await (client as any).initStratum(); + jest.advanceTimersByTime(0); + await Promise.resolve(); + await Promise.resolve(); + firstClientSendSpy.mockClear(); + + const secondSocket = new Socket(); + jest.spyOn(secondSocket, 'on').mockImplementation(() => secondSocket); + secondSocket.end = jest.fn(); + jest.spyOn(secondSocket, 'destroy').mockImplementation(() => secondSocket); + const secondClient = new StratumV1Client( + secondSocket, + stratumV1JobsService, + bitcoinRpcService, + clientService, + clientStatisticsService, + notificationService, + blocksService, + configService, + moduleRef.get(AddressSettingsService) + ); + jest.spyOn(secondClient as any, 'write').mockImplementation((data) => Promise.resolve(true)); + const secondClientSendSpy = jest.spyOn(secondClient as any, 'sendNewMiningJob'); + (secondClient as any).clientSubscription = { userAgent: 'bitaxe' }; + (secondClient as any).clientAuthorization = { + address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', + worker: 'bitaxe4' + }; + (secondClient as any).statistics = {}; + + await (secondClient as any).initStratum(); + jest.advanceTimersByTime(0); + await Promise.resolve(); + await Promise.resolve(); + + expect(secondClientSendSpy).toHaveBeenCalled(); + expect(firstClientSendSpy).not.toHaveBeenCalled(); + expect((StratumV1Client as any).miningJobCache.size).toBe(1); + await secondClient.destroy(); + }); + + it('should skip queued current-job sends when a newer template arrives first', async () => { + jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); + const sendNewMiningJobSpy = jest.spyOn(client as any, 'sendNewMiningJob'); + + (client as any).clientSubscription = { userAgent: 'bitaxe' }; + (client as any).clientAuthorization = { + address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', + worker: 'bitaxe3' + }; + (client as any).statistics = {}; + await (client as any).initStratum(); + jest.advanceTimersByTime(0); + await Promise.resolve(); + await Promise.resolve(); + sendNewMiningJobSpy.mockClear(); + + const oldTemplate = (StratumV1Client as any).latestJobTemplate; + (StratumV1Client as any).queueClientJob(client, oldTemplate); + (StratumV1Client as any).latestJobTemplate = { + ...oldTemplate, + blockData: { + ...oldTemplate.blockData, + id: 'new-template' + } + }; + + jest.advanceTimersByTime(0); + await Promise.resolve(); + await Promise.resolve(); + + expect(sendNewMiningJobSpy).not.toHaveBeenCalled(); + }); + + it('should clear duplicate share tracking immediately when a clear-jobs template is queued', async () => { + jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); + + (client as any).clientSubscription = { userAgent: 'bitaxe' }; + (client as any).clientAuthorization = { + address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', + worker: 'bitaxe3' + }; + (client as any).statistics = {}; + await (client as any).initStratum(); + jest.advanceTimersByTime(0); + await Promise.resolve(); + await Promise.resolve(); + + (client as any).miningSubmissionHashes.add('old-share'); + const latestTemplate = (StratumV1Client as any).latestJobTemplate; + (StratumV1Client as any).queueJobBroadcast({ + ...latestTemplate, + blockData: { + ...latestTemplate.blockData, + id: 'clear-template', + clearJobs: true + } + }); + + expect((client as any).miningSubmissionHashes.size).toBe(0); + }); + + it('should check difficulty through the shared scheduler', async () => { + process.env.STRATUM_DIFFICULTY_CHECK_INTERVAL_MS = '1000'; + process.env.STRATUM_DIFFICULTY_CHECK_BATCH_SIZE = '1'; + jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); + const checkDifficultySpy = jest.spyOn(client as any, 'checkDifficulty').mockResolvedValue(undefined); + + (client as any).clientSubscription = { userAgent: 'bitaxe' }; + (client as any).clientAuthorization = { + address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', + worker: 'bitaxe3' + }; + (client as any).statistics = {}; + await (client as any).initStratum(); + expect((StratumV1Client as any).activeClients.has(client)).toBe(true); + expect((StratumV1Client as any).difficultyCheckInterval).not.toBeNull(); + + jest.advanceTimersByTime(1000); + jest.runOnlyPendingTimers(); + await Promise.resolve(); + await Promise.resolve(); + + expect(checkDifficultySpy).toHaveBeenCalled(); + + delete process.env.STRATUM_DIFFICULTY_CHECK_INTERVAL_MS; + delete process.env.STRATUM_DIFFICULTY_CHECK_BATCH_SIZE; + }); + it('should use the header-only fast path for non-block submissions', async () => { jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); const buildHeaderSpy = jest.spyOn(MiningJob.prototype, 'buildHeaderBuffer'); diff --git a/src/models/StratumV1Client.ts b/src/models/StratumV1Client.ts index c81ed0f..f477030 100644 --- a/src/models/StratumV1Client.ts +++ b/src/models/StratumV1Client.ts @@ -5,7 +5,6 @@ import { validate, ValidationError, ValidatorOptions } from 'class-validator'; import * as crypto from 'crypto'; import { Socket } from 'net'; import { firstValueFrom, Subscription } from 'rxjs'; -import { clearInterval } from 'timers'; import { AddressSettingsService } from '../ORM/address-settings/address-settings.service'; import { BlocksService } from '../ORM/blocks/blocks.service'; @@ -33,13 +32,20 @@ const BLOCKED_USER_AGENT_LOG_INTERVAL_MS = 60 * 1000; const VALIDATION_ERROR_LOG_INTERVAL_MS = 60 * 1000; const DEFAULT_JOB_BROADCAST_BATCH_SIZE = 250; const DEFAULT_HANDSHAKE_TIMEOUT_MS = 30000; +const DEFAULT_DIFFICULTY_CHECK_INTERVAL_MS = 60 * 1000; +const DEFAULT_DIFFICULTY_CHECK_BATCH_SIZE = 1000; -interface JobBroadcastQueueItem { +interface ClientJobQueueItem { client: StratumV1Client; jobTemplate: IJobTemplate; jobTemplateId: string; } +interface PayoutInformation { + address: string; + percent: number; +} + export function effectiveJobDifficulty( jobIdInt: number, currentDiff: number, @@ -59,16 +65,25 @@ export function effectiveJobDifficulty( export class StratumV1Client { private static blockedUserAgentLogState = new Map(); private static validationErrorLogState = new Map(); - private static jobBroadcastQueue: JobBroadcastQueueItem[] = []; + private static activeClients = new Set(); + private static jobBroadcastSubscription: Subscription | null = null; + private static latestJobTemplate: IJobTemplate | null = null; + private static pendingJobTemplate: IJobTemplate | null = null; + private static jobBroadcastIterator: SetIterator | null = null; private static jobBroadcastDraining = false; - private static jobBroadcastDrainOffset = 0; + private static clientJobQueue: ClientJobQueueItem[] = []; + private static clientJobDraining = false; + private static clientJobDrainOffset = 0; + private static difficultyCheckInterval: NodeJS.Timeout | null = null; + private static difficultyCheckIterator: SetIterator | null = null; + private static difficultyCheckDraining = false; + private static miningJobCacheTemplateId: string | null = null; + private static miningJobCache = new Map(); public clientSubscription: SubscriptionMessage; private clientConfiguration: ConfigurationMessage; private clientAuthorization: AuthorizationMessage; private clientSuggestedDifficulty: SuggestDifficulty; - private stratumSubscription: Subscription; - private backgroundWork: NodeJS.Timeout[] = []; private statistics: StratumV1ClientStatistics; private stratumInitialized = false; @@ -142,22 +157,17 @@ export class StratumV1Client { } this.destroyed = true; + StratumV1Client.activeClients.delete(this); + StratumV1Client.stopSharedWorkersIfIdle(); + if (this.clientEntity?.id) { await this.clientService.delete(this.clientEntity.id); } - if (this.stratumSubscription != null) { - this.stratumSubscription.unsubscribe(); - } - if (this.handshakeTimeout != null) { clearTimeout(this.handshakeTimeout); this.handshakeTimeout = null; } - - this.backgroundWork.forEach(work => { - clearInterval(work); - }); } private getRandomHexString() { @@ -443,23 +453,13 @@ export class StratumV1Client { } } - this.stratumSubscription = this.stratumV1JobsService.newMiningJob$.subscribe(async (jobTemplate) => { - try { - if(jobTemplate.blockData.clearJobs){ - this.miningSubmissionHashes.clear(); - } - this.enqueueNewMiningJob(jobTemplate); - } catch (e) { - await this.socket.end(); - console.error(e); - } - }); - - this.backgroundWork.push( - setInterval(async () => { - await this.checkDifficulty(); - }, 60 * 1000) - ); + const sharedJobBroadcastAlreadyRunning = StratumV1Client.jobBroadcastSubscription != null; + StratumV1Client.activeClients.add(this); + StratumV1Client.ensureSharedJobBroadcast(this.stratumV1JobsService); + if (sharedJobBroadcastAlreadyRunning && StratumV1Client.latestJobTemplate != null) { + StratumV1Client.queueClientJob(this, StratumV1Client.latestJobTemplate); + } + StratumV1Client.ensureSharedDifficultyChecks(); // this.backgroundWork.push( // setInterval(async () => { @@ -468,21 +468,54 @@ export class StratumV1Client { // ); } - private enqueueNewMiningJob(jobTemplate: IJobTemplate) { - if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) { + private static ensureSharedJobBroadcast(stratumV1JobsService: StratumV1JobsService) { + if (StratumV1Client.jobBroadcastSubscription != null) { return; } - this.latestQueuedJobTemplateId = jobTemplate.blockData.id; - StratumV1Client.jobBroadcastQueue.push({ - client: this, - jobTemplate, - jobTemplateId: jobTemplate.blockData.id + StratumV1Client.jobBroadcastSubscription = stratumV1JobsService.newMiningJob$.subscribe((jobTemplate) => { + StratumV1Client.latestJobTemplate = jobTemplate; + StratumV1Client.queueJobBroadcast(jobTemplate); }); - StratumV1Client.scheduleJobBroadcastDrain(); } - private static scheduleJobBroadcastDrain() { + private static stopSharedWorkersIfIdle() { + if (StratumV1Client.activeClients.size > 0) { + return; + } + + if (StratumV1Client.jobBroadcastSubscription != null) { + StratumV1Client.jobBroadcastSubscription.unsubscribe(); + StratumV1Client.jobBroadcastSubscription = null; + } + + if (StratumV1Client.difficultyCheckInterval != null) { + clearInterval(StratumV1Client.difficultyCheckInterval); + StratumV1Client.difficultyCheckInterval = null; + } + + StratumV1Client.pendingJobTemplate = null; + StratumV1Client.jobBroadcastIterator = null; + StratumV1Client.jobBroadcastDraining = false; + StratumV1Client.latestJobTemplate = null; + StratumV1Client.clientJobQueue = []; + StratumV1Client.clientJobDraining = false; + StratumV1Client.clientJobDrainOffset = 0; + StratumV1Client.difficultyCheckIterator = null; + StratumV1Client.difficultyCheckDraining = false; + StratumV1Client.miningJobCacheTemplateId = null; + StratumV1Client.miningJobCache.clear(); + } + + private static queueJobBroadcast(jobTemplate: IJobTemplate) { + StratumV1Client.pendingJobTemplate = jobTemplate; + StratumV1Client.jobBroadcastIterator = StratumV1Client.activeClients.values(); + if (jobTemplate.blockData.clearJobs) { + for (const client of StratumV1Client.activeClients) { + client.miningSubmissionHashes.clear(); + } + } + if (StratumV1Client.jobBroadcastDraining) { return; } @@ -491,34 +524,92 @@ export class StratumV1Client { setTimeout(() => StratumV1Client.drainJobBroadcastQueue(), 0); } - private static drainJobBroadcastQueue() { + private static queueClientJob(client: StratumV1Client, jobTemplate: IJobTemplate) { + if (client.connectionClosed || client.socket.destroyed || client.socket.writableEnded) { + return; + } + + client.latestQueuedJobTemplateId = jobTemplate.blockData.id; + StratumV1Client.clientJobQueue.push({ + client, + jobTemplate, + jobTemplateId: jobTemplate.blockData.id + }); + + if (StratumV1Client.clientJobDraining) { + return; + } + + StratumV1Client.clientJobDraining = true; + setTimeout(() => StratumV1Client.drainClientJobQueue(), 0); + } + + private static drainClientJobQueue() { const batchSize = StratumV1Client.getJobBroadcastBatchSize(); let processed = 0; - while (processed < batchSize && StratumV1Client.jobBroadcastDrainOffset < StratumV1Client.jobBroadcastQueue.length) { - const item = StratumV1Client.jobBroadcastQueue[StratumV1Client.jobBroadcastDrainOffset++]; + while (processed < batchSize && StratumV1Client.clientJobDrainOffset < StratumV1Client.clientJobQueue.length) { + const item = StratumV1Client.clientJobQueue[StratumV1Client.clientJobDrainOffset++]; processed++; if (item.client.connectionClosed || item.client.socket.destroyed || item.client.socket.writableEnded + || StratumV1Client.latestJobTemplate?.blockData.id !== item.jobTemplateId || item.client.latestQueuedJobTemplateId !== item.jobTemplateId) { continue; } - item.client.sendNewMiningJob(item.jobTemplate).catch(async (e) => { + item.client.sendNewMiningJob(item.jobTemplate).catch((e) => { item.client.closeSocket(); console.error(e); }); } - if (StratumV1Client.jobBroadcastDrainOffset >= StratumV1Client.jobBroadcastQueue.length) { - StratumV1Client.jobBroadcastQueue = []; - StratumV1Client.jobBroadcastDrainOffset = 0; + if (StratumV1Client.clientJobDrainOffset >= StratumV1Client.clientJobQueue.length) { + StratumV1Client.clientJobQueue = []; + StratumV1Client.clientJobDrainOffset = 0; + StratumV1Client.clientJobDraining = false; + return; + } + + setTimeout(() => StratumV1Client.drainClientJobQueue(), 0); + } + + private static drainJobBroadcastQueue() { + const jobTemplate = StratumV1Client.pendingJobTemplate; + const iterator = StratumV1Client.jobBroadcastIterator; + if (jobTemplate == null || iterator == null) { StratumV1Client.jobBroadcastDraining = false; return; } + const batchSize = StratumV1Client.getJobBroadcastBatchSize(); + let processed = 0; + + while (processed < batchSize) { + const next = iterator.next(); + if (next.done) { + StratumV1Client.jobBroadcastIterator = null; + StratumV1Client.jobBroadcastDraining = false; + return; + } + + const client = next.value; + processed++; + + if (client.connectionClosed || client.socket.destroyed || client.socket.writableEnded) { + continue; + } + + client.latestQueuedJobTemplateId = jobTemplate.blockData.id; + + client.sendNewMiningJob(jobTemplate).catch(async (e) => { + client.closeSocket(); + console.error(e); + }); + } + setTimeout(() => StratumV1Client.drainJobBroadcastQueue(), 0); } @@ -530,9 +621,79 @@ export class StratumV1Client { return DEFAULT_JOB_BROADCAST_BATCH_SIZE; } + private static ensureSharedDifficultyChecks() { + if (StratumV1Client.difficultyCheckInterval != null) { + return; + } + + StratumV1Client.difficultyCheckInterval = setInterval(() => { + StratumV1Client.scheduleDifficultyCheckDrain(); + }, StratumV1Client.getDifficultyCheckIntervalMs()); + } + + private static scheduleDifficultyCheckDrain() { + StratumV1Client.difficultyCheckIterator = StratumV1Client.activeClients.values(); + if (StratumV1Client.difficultyCheckDraining) { + return; + } + + StratumV1Client.difficultyCheckDraining = true; + setTimeout(() => StratumV1Client.drainDifficultyChecks(), 0); + } + + private static drainDifficultyChecks() { + const iterator = StratumV1Client.difficultyCheckIterator; + if (iterator == null) { + StratumV1Client.difficultyCheckDraining = false; + return; + } + + const batchSize = StratumV1Client.getDifficultyCheckBatchSize(); + let processed = 0; + + while (processed < batchSize) { + const next = iterator.next(); + if (next.done) { + StratumV1Client.difficultyCheckIterator = null; + StratumV1Client.difficultyCheckDraining = false; + return; + } + + const client = next.value; + processed++; + + if (client.connectionClosed || client.socket.destroyed || client.socket.writableEnded) { + continue; + } + + client.checkDifficulty().catch((e) => { + client.closeSocket(); + console.error(e); + }); + } + + setTimeout(() => StratumV1Client.drainDifficultyChecks(), 0); + } + + private static getDifficultyCheckIntervalMs() { + const configuredInterval = parseInt(process.env.STRATUM_DIFFICULTY_CHECK_INTERVAL_MS, 10); + if (Number.isFinite(configuredInterval) && configuredInterval > 0) { + return configuredInterval; + } + return DEFAULT_DIFFICULTY_CHECK_INTERVAL_MS; + } + + private static getDifficultyCheckBatchSize() { + const configuredBatchSize = parseInt(process.env.STRATUM_DIFFICULTY_CHECK_BATCH_SIZE, 10); + if (Number.isFinite(configuredBatchSize) && configuredBatchSize > 0) { + return configuredBatchSize; + } + return DEFAULT_DIFFICULTY_CHECK_BATCH_SIZE; + } + private async sendNewMiningJob(jobTemplate: IJobTemplate) { - let payoutInformation= [ + let payoutInformation: PayoutInformation[] = [ { address: this.clientAuthorization.address, percent: 100 } ]; // const devFeeAddress = this.configService.get('DEV_FEE_ADDRESS'); @@ -571,12 +732,14 @@ export class StratumV1Client { throw new Error('Invalid network configuration'); } - const job = new MiningJob( + const poolIdentifier = this.configService.get('POOL_IDENTIFIER') || 'Public-Pool'; + const job = StratumV1Client.getCachedMiningJob( network, this.stratumV1JobsService.getNextId(), payoutInformation, jobTemplate, - this.configService.get('POOL_IDENTIFIER') || 'Public-Pool' + poolIdentifier, + networkConfig ); this.stratumV1JobsService.addJob(job); @@ -592,6 +755,41 @@ export class StratumV1Client { } + private static getCachedMiningJob( + network: bitcoinjs.networks.Network, + jobId: string, + payoutInformation: PayoutInformation[], + jobTemplate: IJobTemplate, + poolIdentifier: string, + networkConfig: string, + ) { + if (StratumV1Client.miningJobCacheTemplateId !== jobTemplate.blockData.id) { + StratumV1Client.miningJobCacheTemplateId = jobTemplate.blockData.id; + StratumV1Client.miningJobCache.clear(); + } + + const cacheKey = JSON.stringify({ + template: jobTemplate.blockData.id, + network: networkConfig, + poolIdentifier, + payoutInformation + }); + let cachedJob = StratumV1Client.miningJobCache.get(cacheKey); + if (cachedJob == null) { + cachedJob = new MiningJob( + network, + jobId, + payoutInformation, + jobTemplate, + poolIdentifier + ); + StratumV1Client.miningJobCache.set(cacheKey, cachedJob); + return cachedJob; + } + + return cachedJob.cloneForJobId(jobId); + } + private async ensureClientEntity() { if (this.clientEntity != null) {