From 9b99c5153a9349331eb8739ee02678fae28502e2 Mon Sep 17 00:00:00 2001 From: Ben Date: Tue, 9 Jun 2026 12:30:03 -0400 Subject: [PATCH] Harden Stratum V1 cleanup --- .env.example | 1 + ecosystem.config.js | 1 + src/models/StratumV1Client.spec.ts | 30 ++++++++++++ src/models/StratumV1Client.ts | 79 ++++++++++++++++++++---------- src/services/stratum-v1.service.ts | 74 ++++++++++++++++++++-------- 5 files changed, 139 insertions(+), 46 deletions(-) diff --git a/.env.example b/.env.example index 69d6bf4..e908fee 100644 --- a/.env.example +++ b/.env.example @@ -28,6 +28,7 @@ DOCKER_LOG_MAX_FILES=5 # Plain TCP Stratum ports accept both SV1 JSON-RPC and SV2 Noise/binary traffic. STRATUM_PORTS=3333,3332,3331,3330 STRATUM_WORKERS=2 +STRATUM_WORKER_MAX_MEMORY_RESTART=4096M STRATUM_MIN_DIFFICULTY=1 STRATUM_SOCKET_TIMEOUT_MS=3600000 STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS=60000 diff --git a/ecosystem.config.js b/ecosystem.config.js index 81d13ac..c4fd413 100644 --- a/ecosystem.config.js +++ b/ecosystem.config.js @@ -42,6 +42,7 @@ module.exports = { script: './dist/main.js', instances: parseInt(process.env.STRATUM_WORKERS || '2', 10), exec_mode: "cluster", + max_memory_restart: process.env.STRATUM_WORKER_MAX_MEMORY_RESTART || '4096M', env: { MASTER: 'false', API_ENABLED: 'false', diff --git a/src/models/StratumV1Client.spec.ts b/src/models/StratumV1Client.spec.ts index 9134888..1720cd6 100644 --- a/src/models/StratumV1Client.spec.ts +++ b/src/models/StratumV1Client.spec.ts @@ -168,6 +168,36 @@ describe('StratumV1Client', () => { expect(socket.on).toHaveBeenCalled(); }); + it('should clean up socket state only once when destroyed repeatedly', async () => { + const unsubscribe = jest.fn(); + const timer = setInterval(() => undefined, 1000); + const removeListenerSpy = jest.spyOn(socket, 'removeListener'); + + (client as any).clientEntity = { + id: '00000000-0000-4000-8000-000000000001', + address: 'tb1qcleanup', + }; + (client as any).stratumSubscription = { unsubscribe }; + (client as any).backgroundWork = [timer]; + (client as any).miningSubmissionHashes.add('submitted-share'); + (client as any).buffer = 'partial-message'; + + await Promise.all([client.destroy(), client.destroy()]); + + expect(redisMessagingService.removeClientPresence).toHaveBeenCalledTimes(1); + expect(redisMessagingService.removeClientPresence).toHaveBeenCalledWith( + '00000000-0000-4000-8000-000000000001', + 'tb1qcleanup', + ); + expect(clientService.delete).toHaveBeenCalledTimes(1); + expect(clientService.delete).toHaveBeenCalledWith('00000000-0000-4000-8000-000000000001'); + expect(unsubscribe).toHaveBeenCalledTimes(1); + expect(removeListenerSpy).toHaveBeenCalledWith('data', expect.any(Function)); + expect((client as any).backgroundWork).toEqual([]); + expect((client as any).miningSubmissionHashes.size).toBe(0); + expect((client as any).buffer).toBe(''); + }); + it('should close socket on invalid JSON', () => { emitMessage('INVALID'); jest.spyOn(socket, 'destroy'); diff --git a/src/models/StratumV1Client.ts b/src/models/StratumV1Client.ts index 7a867ef..0bd9f3f 100644 --- a/src/models/StratumV1Client.ts +++ b/src/models/StratumV1Client.ts @@ -44,6 +44,8 @@ export class StratumV1Client { private clientSuggestedDifficulty: SuggestDifficulty; private stratumSubscription: Subscription; private backgroundWork: NodeJS.Timeout[] = []; + private readonly socketDataHandler: (data: Buffer) => void; + private destroyPromise: Promise | null = null; private statistics: StratumV1ClientStatistics; private stratumInitialized = false; @@ -77,43 +79,68 @@ export class StratumV1Client { private readonly redisMessagingService?: RedisMessagingService ) { - this.socket.on('data', (data: Buffer) => { - this.buffer += data.toString(); - let lines = this.buffer.split('\n'); - this.buffer = lines.pop() || ''; // Save the last part of the data (incomplete line) to the buffer - - (async () => { - for (const m of lines.filter(l => l.length > 0)) { - if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) { - break; - } - try { - await this.handleMessage(m); - } catch (e) { - await this.socket.end(); - console.error(e); - } - } - })(); - }); + this.socketDataHandler = (data: Buffer) => { + void this.handleSocketData(data); + }; + this.socket.on('data', this.socketDataHandler); } - public async destroy() { - - if (this.clientEntity?.id) { - await this.redisMessagingService?.removeClientPresence(this.clientEntity.id, this.clientEntity.address); - await this.clientService.delete(this.clientEntity.id); + public async destroy(): Promise { + if (this.destroyPromise != null) { + return this.destroyPromise; } + this.destroyPromise = this.destroyInternal(); + return this.destroyPromise; + } + + private async destroyInternal(): Promise { + this.connectionClosed = true; + this.socket.removeListener('data', this.socketDataHandler); + this.buffer = ''; + if (this.stratumSubscription != null) { this.stratumSubscription.unsubscribe(); + this.stratumSubscription = null; } - this.backgroundWork.forEach(work => { + for (const work of this.backgroundWork) { clearInterval(work); - }); + } + this.backgroundWork = []; + this.miningSubmissionHashes.clear(); + + if (this.clientEntity?.id) { + const clientId = this.clientEntity.id; + const address = this.clientEntity.address; + this.clientEntity = null; + await this.redisMessagingService?.removeClientPresence(clientId, address); + await this.clientService.delete(clientId); + } + } + + private async handleSocketData(data: Buffer): Promise { + if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) { + return; + } + + this.buffer += data.toString(); + const lines = this.buffer.split('\n'); + this.buffer = lines.pop() || ''; // Save the last part of the data (incomplete line) to the buffer + + for (const m of lines.filter(l => l.length > 0)) { + if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) { + break; + } + try { + await this.handleMessage(m); + } catch (e) { + await this.socket.end(); + console.error(e); + } + } } private getRandomHexString() { diff --git a/src/services/stratum-v1.service.ts b/src/services/stratum-v1.service.ts index 8d53673..758f45b 100644 --- a/src/services/stratum-v1.service.ts +++ b/src/services/stratum-v1.service.ts @@ -125,20 +125,40 @@ export class StratumV1Service implements OnModuleInit { let client: StratumV1Client | StratumV2Client = null; let protocol: 'v1' | 'v2' | null = null; + let cleanedUp = false; // Unified cleanup function const cleanup = async (reason: string) => { - if (client != null && (protocol === 'v2' || (client as StratumV1Client).extraNonceAndSessionId != null)) { - await client.destroy(); - if (reason == 'Error') { - this.errorClosure++; - } else { - this.normalClosure++; - } + if (cleanedUp) { + return; } - if (!socket.destroyed) { - socket.end(); - socket.destroy(); + cleanedUp = true; + + const currentClient = client; + client = null; + + try { + if (currentClient != null) { + const initializedClient = protocol === 'v2' + || (currentClient as StratumV1Client).extraNonceAndSessionId != null; + await currentClient.destroy(); + if (initializedClient) { + if (reason == 'Error') { + this.errorClosure++; + } else { + this.normalClosure++; + } + } + } + } finally { + socket.removeAllListeners('close'); + socket.removeAllListeners('timeout'); + socket.removeAllListeners('error'); + socket.removeAllListeners('data'); + if (!socket.destroyed) { + socket.end(); + socket.destroy(); + } } }; @@ -241,19 +261,33 @@ export class StratumV1Service implements OnModuleInit { socket.setKeepAlive(true, this.getTcpKeepAliveInitialDelayMs()); const client = this.createV1Client(socket); + let cleanedUp = false; const cleanup = async (reason: string) => { - if (client.extraNonceAndSessionId != null) { - await client.destroy(); - if (reason === 'Error') { - this.errorClosure++; - } else { - this.normalClosure++; - } + if (cleanedUp) { + return; } - if (!socket.destroyed) { - socket.end(); - socket.destroy(); + cleanedUp = true; + + try { + const initializedClient = client.extraNonceAndSessionId != null; + await client.destroy(); + if (initializedClient) { + if (reason === 'Error') { + this.errorClosure++; + } else { + this.normalClosure++; + } + } + } finally { + socket.removeAllListeners('close'); + socket.removeAllListeners('timeout'); + socket.removeAllListeners('error'); + socket.removeAllListeners('data'); + if (!socket.destroyed) { + socket.end(); + socket.destroy(); + } } };