From d6d372711969014447568fd0a079073ec5acc229 Mon Sep 17 00:00:00 2001 From: Ben Date: Wed, 6 May 2026 20:38:42 -0400 Subject: [PATCH] improve stratum connection and broadcast performance --- src/models/StratumV1Client.spec.ts | 64 +++++++++++++++++ src/models/StratumV1Client.ts | 111 ++++++++++++++++++++++++++++- src/services/stratum-v1.service.ts | 28 ++++++-- 3 files changed, 198 insertions(+), 5 deletions(-) diff --git a/src/models/StratumV1Client.spec.ts b/src/models/StratumV1Client.spec.ts index ab81e8f..4d9ef98 100644 --- a/src/models/StratumV1Client.spec.ts +++ b/src/models/StratumV1Client.spec.ts @@ -131,6 +131,9 @@ describe('StratumV1Client', () => { }); (StratumV1Client as any).blockedUserAgentLogState.clear(); (StratumV1Client as any).validationErrorLogState.clear(); + (StratumV1Client as any).jobBroadcastQueue = []; + (StratumV1Client as any).jobBroadcastDraining = false; + (StratumV1Client as any).jobBroadcastDrainOffset = 0; bitcoinRpcService = { newBlockTemplate$: newBlockEmitter.asObservable(), @@ -299,6 +302,41 @@ describe('StratumV1Client', () => { await secondClient.destroy(); }); + it('should close clients that do not complete the stratum handshake', async () => { + (configService.get as jest.Mock).mockImplementation((key: string) => { + switch (key) { + case 'STRATUM_HANDSHAKE_TIMEOUT_MS': + return '1000'; + case 'DEV_FEE_ADDRESS': + return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4'; + case 'NETWORK': + return 'testnet'; + } + return null; + }); + + const slowSocket = new Socket(); + jest.spyOn(slowSocket, 'on').mockImplementation(() => slowSocket); + jest.spyOn(slowSocket, 'destroy').mockImplementation(() => slowSocket); + slowSocket.end = jest.fn(); + const slowClient = new StratumV1Client( + slowSocket, + stratumV1JobsService, + bitcoinRpcService, + clientService, + clientStatisticsService, + notificationService, + blocksService, + configService, + moduleRef.get(AddressSettingsService) + ); + + jest.advanceTimersByTime(1000); + + expect(slowSocket.destroy).toHaveBeenCalled(); + await slowClient.destroy(); + }); + it('should respond to mining.configure', async () => { @@ -394,6 +432,32 @@ describe('StratumV1Client', () => { }); + it('should queue mining job broadcasts instead of building jobs in the observable callback', 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' }; + (client as any).clientAuthorization = { + address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', + worker: 'bitaxe3' + }; + (client as any).statistics = {}; + await (client as any).initStratum(); + await Promise.resolve(); + await Promise.resolve(); + + expect(enqueueNewMiningJobSpy).toHaveBeenCalled(); + + jest.advanceTimersByTime(0); + await Promise.resolve(); + await Promise.resolve(); + + expect(sendNewMiningJobSpy).toHaveBeenCalled(); + }); + 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 11925cf..c81ed0f 100644 --- a/src/models/StratumV1Client.ts +++ b/src/models/StratumV1Client.ts @@ -31,6 +31,14 @@ import { StratumV1ClientStatistics } from './StratumV1ClientStatistics'; const TRUE_DIFF_ONE = 2.695953529101131e67; 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; + +interface JobBroadcastQueueItem { + client: StratumV1Client; + jobTemplate: IJobTemplate; + jobTemplateId: string; +} export function effectiveJobDifficulty( jobIdInt: number, @@ -51,6 +59,9 @@ export function effectiveJobDifficulty( export class StratumV1Client { private static blockedUserAgentLogState = new Map(); private static validationErrorLogState = new Map(); + private static jobBroadcastQueue: JobBroadcastQueueItem[] = []; + private static jobBroadcastDraining = false; + private static jobBroadcastDrainOffset = 0; public clientSubscription: SubscriptionMessage; private clientConfiguration: ConfigurationMessage; @@ -76,6 +87,9 @@ export class StratumV1Client { private buffer: string = ''; private connectionClosed = false; + private destroyed = false; + private handshakeTimeout: NodeJS.Timeout | null; + private latestQueuedJobTemplateId: string; private miningSubmissionHashes = new Set() @@ -91,6 +105,12 @@ export class StratumV1Client { private readonly addressSettingsService: AddressSettingsService ) { + this.handshakeTimeout = setTimeout(() => { + if (!this.stratumInitialized) { + this.closeSocket(); + } + }, this.getHandshakeTimeoutMs()); + this.socket.on('data', (data: Buffer) => { this.buffer += data.toString(); let lines = this.buffer.split('\n'); @@ -115,6 +135,12 @@ export class StratumV1Client { } public async destroy() { + this.connectionClosed = true; + + if (this.destroyed) { + return; + } + this.destroyed = true; if (this.clientEntity?.id) { await this.clientService.delete(this.clientEntity.id); @@ -124,6 +150,11 @@ export class StratumV1Client { this.stratumSubscription.unsubscribe(); } + if (this.handshakeTimeout != null) { + clearTimeout(this.handshakeTimeout); + this.handshakeTimeout = null; + } + this.backgroundWork.forEach(work => { clearInterval(work); }); @@ -385,6 +416,10 @@ export class StratumV1Client { private async initStratum() { this.stratumInitialized = true; + if (this.handshakeTimeout != null) { + clearTimeout(this.handshakeTimeout); + this.handshakeTimeout = null; + } if (this.isBlockedUserAgent(this.clientSubscription.userAgent)) { this.logBlockedUserAgent(this.clientSubscription.userAgent); @@ -413,7 +448,7 @@ export class StratumV1Client { if(jobTemplate.blockData.clearJobs){ this.miningSubmissionHashes.clear(); } - await this.sendNewMiningJob(jobTemplate); + this.enqueueNewMiningJob(jobTemplate); } catch (e) { await this.socket.end(); console.error(e); @@ -433,6 +468,68 @@ export class StratumV1Client { // ); } + private enqueueNewMiningJob(jobTemplate: IJobTemplate) { + if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) { + return; + } + + this.latestQueuedJobTemplateId = jobTemplate.blockData.id; + StratumV1Client.jobBroadcastQueue.push({ + client: this, + jobTemplate, + jobTemplateId: jobTemplate.blockData.id + }); + StratumV1Client.scheduleJobBroadcastDrain(); + } + + private static scheduleJobBroadcastDrain() { + if (StratumV1Client.jobBroadcastDraining) { + return; + } + + StratumV1Client.jobBroadcastDraining = true; + setTimeout(() => StratumV1Client.drainJobBroadcastQueue(), 0); + } + + private static drainJobBroadcastQueue() { + const batchSize = StratumV1Client.getJobBroadcastBatchSize(); + let processed = 0; + + while (processed < batchSize && StratumV1Client.jobBroadcastDrainOffset < StratumV1Client.jobBroadcastQueue.length) { + const item = StratumV1Client.jobBroadcastQueue[StratumV1Client.jobBroadcastDrainOffset++]; + processed++; + + if (item.client.connectionClosed + || item.client.socket.destroyed + || item.client.socket.writableEnded + || item.client.latestQueuedJobTemplateId !== item.jobTemplateId) { + continue; + } + + item.client.sendNewMiningJob(item.jobTemplate).catch(async (e) => { + item.client.closeSocket(); + console.error(e); + }); + } + + if (StratumV1Client.jobBroadcastDrainOffset >= StratumV1Client.jobBroadcastQueue.length) { + StratumV1Client.jobBroadcastQueue = []; + StratumV1Client.jobBroadcastDrainOffset = 0; + StratumV1Client.jobBroadcastDraining = false; + return; + } + + setTimeout(() => StratumV1Client.drainJobBroadcastQueue(), 0); + } + + private static getJobBroadcastBatchSize() { + const configuredBatchSize = parseInt(process.env.STRATUM_JOB_BROADCAST_BATCH_SIZE, 10); + if (Number.isFinite(configuredBatchSize) && configuredBatchSize > 0) { + return configuredBatchSize; + } + return DEFAULT_JOB_BROADCAST_BATCH_SIZE; + } + private async sendNewMiningJob(jobTemplate: IJobTemplate) { let payoutInformation= [ @@ -814,8 +911,20 @@ export class StratumV1Client { return ` sample=${values.join(',')}`; } + private getHandshakeTimeoutMs() { + const configuredTimeout = parseInt(this.configService.get('STRATUM_HANDSHAKE_TIMEOUT_MS') ?? process.env.STRATUM_HANDSHAKE_TIMEOUT_MS, 10); + if (Number.isFinite(configuredTimeout) && configuredTimeout > 0) { + return configuredTimeout; + } + return DEFAULT_HANDSHAKE_TIMEOUT_MS; + } + private closeSocket() { this.connectionClosed = true; + if (this.handshakeTimeout != null) { + clearTimeout(this.handshakeTimeout); + this.handshakeTimeout = null; + } if (!this.socket.destroyed) { this.socket.destroy(); } diff --git a/src/services/stratum-v1.service.ts b/src/services/stratum-v1.service.ts index 442c6ae..28dda1c 100644 --- a/src/services/stratum-v1.service.ts +++ b/src/services/stratum-v1.service.ts @@ -70,6 +70,8 @@ export class StratumV1Service implements OnModuleInit { private startSocketServer(port: number) { const server = new Server(async (socket: Socket) => { + socket.setNoDelay(true); + socket.setKeepAlive(true, 60000); // Set 15-minute timeout socket.setTimeout(1000 * 60 * 15); @@ -86,9 +88,17 @@ export class StratumV1Service implements OnModuleInit { ); // Unified cleanup function + let cleanedUp = false; const cleanup = async (reason: string) => { - if (client.extraNonceAndSessionId != null) { - await client.destroy(); + if (cleanedUp) { + return; + } + cleanedUp = true; + + const initializedClient = client.extraNonceAndSessionId != null; + await client.destroy(); + + if (initializedClient) { if (reason == 'Error') { this.errorClosure++; } else { @@ -149,6 +159,8 @@ export class StratumV1Service implements OnModuleInit { }; const server = createServer(tlsOptions, async (socket: TLSSocket) => { + socket.setNoDelay(true); + socket.setKeepAlive(true, 60000); // Set 15-minute timeout socket.setTimeout(1000 * 60 * 15); @@ -164,9 +176,17 @@ export class StratumV1Service implements OnModuleInit { this.addressSettingsService ); + let cleanedUp = false; const cleanup = async (reason: string) => { - if (client.extraNonceAndSessionId != null) { - await client.destroy(); + if (cleanedUp) { + return; + } + cleanedUp = true; + + const initializedClient = client.extraNonceAndSessionId != null; + await client.destroy(); + + if (initializedClient) { if (reason === 'Error') { this.errorClosure++; } else {