From 8709de3d83a1ee3d63da77f0d39efe395c86cf8e Mon Sep 17 00:00:00 2001 From: Ben Date: Wed, 6 May 2026 20:26:00 -0400 Subject: [PATCH] reduce stratum rejection overhead --- src/models/StratumV1Client.spec.ts | 44 ++++++++++++++ src/models/StratumV1Client.ts | 55 ++++++++++++++++- src/services/stratum-v1.service.spec.ts | 78 +++++++++++++++++++++++++ src/services/stratum-v1.service.ts | 4 +- 4 files changed, 177 insertions(+), 4 deletions(-) create mode 100644 src/services/stratum-v1.service.spec.ts diff --git a/src/models/StratumV1Client.spec.ts b/src/models/StratumV1Client.spec.ts index 9e9c720..ab81e8f 100644 --- a/src/models/StratumV1Client.spec.ts +++ b/src/models/StratumV1Client.spec.ts @@ -130,6 +130,7 @@ describe('StratumV1Client', () => { return null; }); (StratumV1Client as any).blockedUserAgentLogState.clear(); + (StratumV1Client as any).validationErrorLogState.clear(); bitcoinRpcService = { newBlockTemplate$: newBlockEmitter.asObservable(), @@ -540,6 +541,49 @@ describe('StratumV1Client', () => { await new Promise((r) => setTimeout(r, 100)); expect((client as any).write).lastCalledWith(expect.stringContaining(`"error":[20,"Mining Submit validation error"`)); + expect(socket.destroy).toHaveBeenCalled(); + expect(consoleWarnSpy).toHaveBeenCalledWith(expect.stringContaining('Mining Submit validation error: extraNonce2:isLength')); + }); + + it('should throttle repeated mining submit validation logs', async () => { + jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); + + emitMessage(MockRecording1.MINING_SUBSCRIBE); + emitMessage(MockRecording1.MINING_AUTHORIZE); + await new Promise((r) => setTimeout(r, 100)); + + emitMessage(`{"id": 5, "method": "mining.submit", "params": ["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.bitaxe3", "1", "c7080000", "64b3f3ec", "ed460d91", "00002000"]}`); + await new Promise((r) => setTimeout(r, 100)); + + const secondSocket = new Socket(); + jest.spyOn(secondSocket, 'on').mockImplementation((event: string, listener: (...args: any[]) => void) => { + socketEmitter = listener; + return 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)); + jest.spyOn(secondClient as any, 'getRandomHexString').mockReturnValue(MockRecording1.EXTRA_NONCE); + + socketEmitter(Buffer.from(`${MockRecording1.MINING_SUBSCRIBE}\n`)); + socketEmitter(Buffer.from(`${MockRecording1.MINING_AUTHORIZE}\n`)); + await new Promise((r) => setTimeout(r, 100)); + socketEmitter(Buffer.from(`{"id": 5, "method": "mining.submit", "params": ["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.bitaxe3", "1", "c7080000", "64b3f3ec", "ed460d91", "00002000"]}\n`)); + await new Promise((r) => setTimeout(r, 100)); + + expect(consoleWarnSpy.mock.calls.filter(call => call[0]?.startsWith('Mining Submit validation error'))).toHaveLength(1); + await secondClient.destroy(); }); it('should close socket when a submit arrives before stratum is initialized', async () => { diff --git a/src/models/StratumV1Client.ts b/src/models/StratumV1Client.ts index 33e7b21..11925cf 100644 --- a/src/models/StratumV1Client.ts +++ b/src/models/StratumV1Client.ts @@ -1,7 +1,7 @@ import { ConfigService } from '@nestjs/config'; import * as bitcoinjs from 'bitcoinjs-lib'; import { plainToInstance } from 'class-transformer'; -import { validate, ValidatorOptions } from 'class-validator'; +import { validate, ValidationError, ValidatorOptions } from 'class-validator'; import * as crypto from 'crypto'; import { Socket } from 'net'; import { firstValueFrom, Subscription } from 'rxjs'; @@ -30,6 +30,7 @@ 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; export function effectiveJobDifficulty( jobIdInt: number, @@ -49,6 +50,7 @@ export function effectiveJobDifficulty( export class StratumV1Client { private static blockedUserAgentLogState = new Map(); + private static validationErrorLogState = new Map(); public clientSubscription: SubscriptionMessage; private clientConfiguration: ConfigurationMessage; @@ -348,17 +350,18 @@ export class StratumV1Client { } else { - console.log('Mining Submit validation error'); + this.logValidationError('Mining Submit validation error', errors); const err = new StratumErrorMessage( miningSubmitMessage.id, eStratumErrorCode.OtherUnknown, 'Mining Submit validation error', errors).response(); - console.error(err); const success = await this.write(err); if (!success) { return; } + this.closeSocket(); + return; } break; } @@ -765,6 +768,52 @@ export class StratumV1Client { }); } + private logValidationError(label: string, errors: ValidationError[]) { + const now = Date.now(); + const signature = this.getValidationErrorSignature(errors); + const sample = this.getValidationErrorSample(errors); + const key = `${label}:${signature}`; + const logState = StratumV1Client.validationErrorLogState.get(key); + + if (logState != null && now < logState.nextLogAt) { + logState.suppressed += 1; + return; + } + + const suppressed = logState?.suppressed ?? 0; + const suffix = suppressed > 0 ? ` (${suppressed} similar validation errors suppressed)` : ''; + console.warn(`${label}: ${signature}${sample}${suffix}`); + StratumV1Client.validationErrorLogState.set(key, { + nextLogAt: now + VALIDATION_ERROR_LOG_INTERVAL_MS, + suppressed: 0, + sample + }); + } + + private getValidationErrorSignature(errors: ValidationError[]): string { + if (errors.length === 0) { + return 'unknown'; + } + + return errors.map(error => { + const constraints = Object.keys(error.constraints ?? {}).sort().join('|') || 'invalid'; + return `${error.property}:${constraints}`; + }).join(';'); + } + + private getValidationErrorSample(errors: ValidationError[]): string { + const values = errors + .map(error => error.value) + .filter(value => value != null) + .map(value => String(value).replace(/[\r\n]/g, '').slice(0, 64)); + + if (values.length === 0) { + return ''; + } + + return ` sample=${values.join(',')}`; + } + private closeSocket() { this.connectionClosed = true; if (!this.socket.destroyed) { diff --git a/src/services/stratum-v1.service.spec.ts b/src/services/stratum-v1.service.spec.ts new file mode 100644 index 0000000..41c47e1 --- /dev/null +++ b/src/services/stratum-v1.service.spec.ts @@ -0,0 +1,78 @@ +import { StratumV1Service } from './stratum-v1.service'; + +describe('StratumV1Service', () => { + const originalMaster = process.env.MASTER; + const originalStratumPorts = process.env.STRATUM_PORTS; + const originalStratumSecure = process.env.STRATUM_SECURE; + const originalSecureStratumPorts = process.env.SECURE_STRATUM_PORTS; + + let service: StratumV1Service; + let clientService; + let consoleLogSpy: jest.SpyInstance; + + beforeEach(() => { + jest.useFakeTimers(); + clientService = { + deleteAll: jest.fn().mockResolvedValue(undefined) + }; + service = new StratumV1Service( + {} as any, + clientService, + {} as any, + {} as any, + {} as any, + {} as any, + {} as any, + {} as any + ); + consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined); + }); + + afterEach(() => { + restoreEnv('MASTER', originalMaster); + restoreEnv('STRATUM_PORTS', originalStratumPorts); + restoreEnv('STRATUM_SECURE', originalStratumSecure); + restoreEnv('SECURE_STRATUM_PORTS', originalSecureStratumPorts); + consoleLogSpy.mockRestore(); + jest.useRealTimers(); + }); + + it('should skip Stratum listeners in the master process', async () => { + process.env.MASTER = 'true'; + const startSocketServerSpy = jest.spyOn(service as any, 'startSocketServer'); + const startSecureSocketServerSpy = jest.spyOn(service as any, 'startSecureSocketServer'); + + await service.onModuleInit(); + jest.runOnlyPendingTimers(); + + expect(clientService.deleteAll).toHaveBeenCalled(); + expect(startSocketServerSpy).not.toHaveBeenCalled(); + expect(startSecureSocketServerSpy).not.toHaveBeenCalled(); + expect(consoleLogSpy).toHaveBeenCalledWith('Master process skipping Stratum socket listeners'); + }); + + it('should start Stratum listeners in worker processes', async () => { + process.env.MASTER = 'false'; + process.env.STRATUM_PORTS = '3333,3334'; + process.env.STRATUM_SECURE = 'true'; + process.env.SECURE_STRATUM_PORTS = '4333'; + const startSocketServerSpy = jest.spyOn(service as any, 'startSocketServer').mockImplementation(() => undefined); + const startSecureSocketServerSpy = jest.spyOn(service as any, 'startSecureSocketServer').mockImplementation(() => undefined); + + await service.onModuleInit(); + jest.advanceTimersByTime(10000); + + expect(clientService.deleteAll).not.toHaveBeenCalled(); + expect(startSocketServerSpy).toHaveBeenCalledWith(3333); + expect(startSocketServerSpy).toHaveBeenCalledWith(3334); + expect(startSecureSocketServerSpy).toHaveBeenCalledWith(4333); + }); + + function restoreEnv(key: string, value: string | undefined) { + if (value == null) { + delete process.env[key]; + return; + } + process.env[key] = value; + } +}); diff --git a/src/services/stratum-v1.service.ts b/src/services/stratum-v1.service.ts index aec393b..442c6ae 100644 --- a/src/services/stratum-v1.service.ts +++ b/src/services/stratum-v1.service.ts @@ -42,6 +42,8 @@ export class StratumV1Service implements OnModuleInit { if (process.env.MASTER == 'true') { await this.clientService.deleteAll(); + console.log('Master process skipping Stratum socket listeners'); + return; } // wait for all the other processes to init for an even connection distribution @@ -208,4 +210,4 @@ export class StratumV1Service implements OnModuleInit { } -} \ No newline at end of file +}