From 07a2b173d036f8c3db3735e778b7939fc4a34ad2 Mon Sep 17 00:00:00 2001 From: Ben Date: Wed, 5 Aug 2026 16:34:43 -0400 Subject: [PATCH] Harden Stratum V1 share validation --- .env.example | 1 + src/models/StratumV1Client.spec.ts | 95 +++++++++++++++++++- src/models/StratumV1Client.ts | 95 ++++++++++++++------ src/models/StratumV1ClientStatistics.spec.ts | 16 ++++ src/models/StratumV1ClientStatistics.ts | 28 +++--- 5 files changed, 187 insertions(+), 48 deletions(-) diff --git a/.env.example b/.env.example index 27412f7..d8e4adb 100644 --- a/.env.example +++ b/.env.example @@ -72,6 +72,7 @@ STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS=60000 # Disconnect slow/non-reading SV1 clients before their per-socket write queue # can grow without bound. Includes bytes already buffered by Node.js. STRATUM_MAX_SOCKET_BUFFER_BYTES=262144 +STRATUM_MAX_INBOUND_LINE_BYTES=65536 # Immediately publish a consensus-valid, subsidy-only solo job before the full # transaction template. SV2 uses the same activation to select its pre-staged # future job; the full job then follows on the active tip. diff --git a/src/models/StratumV1Client.spec.ts b/src/models/StratumV1Client.spec.ts index 3ad3307..77a9f33 100644 --- a/src/models/StratumV1Client.spec.ts +++ b/src/models/StratumV1Client.spec.ts @@ -169,6 +169,59 @@ describe('StratumV1Client', () => { expect(socket.on).toHaveBeenCalled(); }); + it('disconnects when an unterminated inbound message exceeds the configured limit', async () => { + (configService.get as jest.Mock).mockImplementation((key: string) => { + if (key === 'STRATUM_MAX_INBOUND_LINE_BYTES') return '16'; + if (key === 'NETWORK') return 'testnet'; + return null; + }); + client = new StratumV1Client( + socket, + stratumV1JobsService, + bitcoinRpcService, + clientService, + notificationService, + blocksService, + configService, + addressSettings, + shareAccountingService as any, + redisMessagingService as any, + ); + + socketEmitter(Buffer.from('x'.repeat(17))); + await Promise.resolve(); + + expect(socket.destroy).toHaveBeenCalled(); + expect((client as any).buffer).toBe(''); + }); + + it('accepts multiple complete messages when each line is within the inbound limit', async () => { + (configService.get as jest.Mock).mockImplementation((key: string) => { + if (key === 'STRATUM_MAX_INBOUND_LINE_BYTES') return '64'; + if (key === 'NETWORK') return 'testnet'; + return null; + }); + client = new StratumV1Client( + socket, + stratumV1JobsService, + bitcoinRpcService, + clientService, + notificationService, + blocksService, + configService, + addressSettings, + shareAccountingService as any, + redisMessagingService as any, + ); + jest.spyOn(client as any, 'handleMessage').mockResolvedValue(undefined); + + socketEmitter(Buffer.from('{"id":1}\n{"id":2}\n')); + await Promise.resolve(); + + expect((client as any).handleMessage).toHaveBeenCalledTimes(2); + expect(socket.destroy).not.toHaveBeenCalled(); + }); + it('should clean up socket state only once when destroyed repeatedly', async () => { const timer = setInterval(() => undefined, 1000); const removeListenerSpy = jest.spyOn(socket, 'removeListener'); @@ -805,6 +858,7 @@ describe('StratumV1Client', () => { await new Promise((r) => setTimeout(r, 100)); expect((client as any).isDuplicateSubmission('old-tip-share')).toBe(false); + expect((client as any).rememberSubmission('old-tip-share', '1')).toBe(true); const nextTip = { ...MockRecording1.BLOCK_TEMPLATE, previousblockhash: '11'.repeat(32), @@ -817,7 +871,9 @@ describe('StratumV1Client', () => { expect((client as any).isDuplicateSubmission('old-tip-share')).toBe(true); }); - it('should bound duplicate tracking and expire entries by TTL', () => { + it('should preserve live accepted-share dedup entries when capacity is reached', () => { + const getSubmissionContext = jest.spyOn(stratumV1JobsService, 'getSubmissionContext') + .mockReturnValue({} as any); (configService.get as jest.Mock).mockImplementation((key: string) => { if (key === 'STRATUM_SUBMISSION_DEDUP_TTL_MS') return '1000'; if (key === 'STRATUM_SUBMISSION_DEDUP_MAX_ENTRIES') return '2'; @@ -826,13 +882,18 @@ describe('StratumV1Client', () => { }); expect((client as any).isDuplicateSubmission('one')).toBe(false); + expect((client as any).rememberSubmission('one', '1')).toBe(true); expect((client as any).isDuplicateSubmission('two')).toBe(false); - expect((client as any).isDuplicateSubmission('three')).toBe(false); - expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['two', 'three']); + expect((client as any).rememberSubmission('two', '1')).toBe(true); + expect((client as any).rememberSubmission('three', '1')).toBe(false); + expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['one', 'two']); jest.advanceTimersByTime(1001); + expect((client as any).isDuplicateSubmission('two')).toBe(true); + getSubmissionContext.mockReturnValue(null); expect((client as any).isDuplicateSubmission('two')).toBe(false); - expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['two']); + expect((client as any).rememberSubmission('three', '1')).toBe(true); + expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['three']); }); it('should hash and reject stale non-block shares without accounting them', async () => { @@ -981,6 +1042,32 @@ describe('StratumV1Client', () => { expect((client as any).write).lastCalledWith(`{"id":5,"result":null,"error":[23,"Difficulty too low",""]}\n`); expect(await clientService.connectedClientCount()).toBe(1); expect(shareAccountingService.recordAcceptedShare).not.toHaveBeenCalled(); + expect((client as any).miningSubmissionHashes.size).toBe(0); + }); + + it.each([ + { label: 'before the advertised job time', offsetSeconds: -1 }, + { label: 'more than two hours in the future', offsetSeconds: (2 * 60 * 60) + 1 }, + ])('rejects ntime $label before hashing or accounting', async ({ offsetSeconds }) => { + jest.spyOn(client as any, 'write').mockResolvedValue(true); + const calculateDifficultySpy = jest.spyOn(client as any, 'calculateDifficulty'); + + emitMessage(MockRecording1.MINING_SUBSCRIBE); + emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`); + emitMessage(MockRecording1.MINING_AUTHORIZE); + await new Promise((resolve) => setTimeout(resolve, 100)); + + const submission = JSON.parse(MockRecording1.MINING_SUBMIT); + const baseTime = parseInt(MockRecording1.TIME, 16); + submission.params[3] = (baseTime + offsetSeconds).toString(16).padStart(8, '0'); + emitMessage(JSON.stringify(submission)); + await new Promise((resolve) => setTimeout(resolve, 100)); + + expect((client as any).write).lastCalledWith( + `{"id":5,"result":null,"error":[20,"Invalid ntime",""]}\n`, + ); + expect(calculateDifficultySpy).not.toHaveBeenCalled(); + expect(shareAccountingService.recordAcceptedShare).not.toHaveBeenCalled(); }); it('rejects version rolling masks outside the negotiated BIP320 range', async () => { diff --git a/src/models/StratumV1Client.ts b/src/models/StratumV1Client.ts index 581214c..cd9f0a3 100644 --- a/src/models/StratumV1Client.ts +++ b/src/models/StratumV1Client.ts @@ -46,6 +46,8 @@ const DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS = 60 * 1000; const DEFAULT_SUBMISSION_DEDUP_TTL_MS = 5 * 60 * 1000; const DEFAULT_SUBMISSION_DEDUP_MAX_ENTRIES = 10_000; const DEFAULT_MAX_SOCKET_BUFFER_BYTES = 256 * 1024; +const DEFAULT_MAX_INBOUND_LINE_BYTES = 64 * 1024; +const MAX_NTIME_FUTURE_SECONDS = 2 * 60 * 60; const VERSION_ROLLING_MASK = 0x1fffe000; export interface MiningJobBroadcastResult { @@ -55,6 +57,11 @@ export interface MiningJobBroadcastResult { preStaged?: boolean; } +interface SubmissionDedupEntry { + expiresAt: number; + jobId: string; +} + export class StratumV1Client { private static blockedUserAgentLogState = new Map(); private static validationErrorLogState = new Map(); @@ -90,8 +97,9 @@ export class StratumV1Client { private lastHashRatePersistedAt = 0; private readonly network: bitcoinjs.Network; private readonly maxSocketBufferBytes: number; + private readonly maxInboundLineBytes: number; - private miningSubmissionHashes = new Map(); + private miningSubmissionHashes = new Map(); constructor( public readonly socket: Socket, @@ -115,6 +123,7 @@ export class StratumV1Client { this.socket.on('data', this.socketDataHandler); this.network = this.getNetwork(); this.maxSocketBufferBytes = this.readMaxSocketBufferBytes(); + this.maxInboundLineBytes = this.readMaxInboundLineBytes(); } @@ -150,9 +159,15 @@ export class StratumV1Client { 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 + const lines = `${this.buffer}${data.toString()}`.split('\n'); + const incompleteLine = lines.pop() || ''; + if (Buffer.byteLength(incompleteLine) > this.maxInboundLineBytes + || lines.some(line => Buffer.byteLength(line) > this.maxInboundLineBytes)) { + this.buffer = ''; + this.closeSocket(); + return; + } + this.buffer = incompleteLine; for (const m of lines.filter(l => l.length > 0)) { if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) { @@ -679,6 +694,15 @@ export class StratumV1Client { ); return false; } + const maximumNtime = Math.floor(Date.now() / 1000) + MAX_NTIME_FUTURE_SECONDS; + if (timestamp < jobTemplate.block.timestamp || timestamp > maximumNtime) { + await this.writeSubmissionError( + submission, + eStratumErrorCode.OtherUnknown, + 'Invalid ntime', + ); + return false; + } // The optional BIP310 field contains replacement bits, not an XOR // delta. A legacy five-field submission leaves the advertised version // unchanged, including any bits already set inside the rolling mask. @@ -757,6 +781,15 @@ export class StratumV1Client { const creditedDifficulty = isBlockCandidate && !meetsSessionTarget ? Math.min(this.sessionDifficulty, submissionDifficulty) : this.sessionDifficulty; + if (!this.rememberSubmission(submissionHash, job.jobId)) { + await this.writeSubmissionError( + submission, + eStratumErrorCode.OtherUnknown, + 'Submission dedup capacity exceeded', + ); + this.closeSocket(); + return false; + } let blockSubmissionResult: string = null; if (status === 'stale') { @@ -1052,33 +1085,31 @@ export class StratumV1Client { private isDuplicateSubmission(submissionHash: string): boolean { const now = Date.now(); - const existingExpiry = this.miningSubmissionHashes.get(submissionHash); - if (existingExpiry != null && existingExpiry > now) { - return true; - } - if (existingExpiry != null) { - this.miningSubmissionHashes.delete(submissionHash); - } + this.pruneExpiredSubmissions(now); + return this.miningSubmissionHashes.has(submissionHash); + } - for (const [hash, expiresAt] of this.miningSubmissionHashes) { - if (expiresAt <= now) { + private rememberSubmission(submissionHash: string, jobId: string): boolean { + const now = Date.now(); + this.pruneExpiredSubmissions(now); + if (this.miningSubmissionHashes.has(submissionHash) + || this.miningSubmissionHashes.size >= this.getSubmissionDedupMaxEntries()) { + return false; + } + this.miningSubmissionHashes.set(submissionHash, { + expiresAt: now + this.getSubmissionDedupTtlMs(), + jobId, + }); + return true; + } + + private pruneExpiredSubmissions(now: number): void { + for (const [hash, entry] of this.miningSubmissionHashes) { + if (entry.expiresAt <= now + && this.stratumV1JobsService.getSubmissionContext(entry.jobId) == null) { this.miningSubmissionHashes.delete(hash); } } - - this.miningSubmissionHashes.set( - submissionHash, - now + this.getSubmissionDedupTtlMs(), - ); - const maxEntries = this.getSubmissionDedupMaxEntries(); - while (this.miningSubmissionHashes.size > maxEntries) { - const oldestHash = this.miningSubmissionHashes.keys().next().value; - if (oldestHash == null) { - break; - } - this.miningSubmissionHashes.delete(oldestHash); - } - return false; } private getSubmissionDedupTtlMs(): number { @@ -1111,6 +1142,16 @@ export class StratumV1Client { : DEFAULT_MAX_SOCKET_BUFFER_BYTES; } + private readMaxInboundLineBytes(): number { + const configured = Number( + this.configService.get('STRATUM_MAX_INBOUND_LINE_BYTES') + ?? process.env.STRATUM_MAX_INBOUND_LINE_BYTES, + ); + return Number.isSafeInteger(configured) && configured > 0 + ? configured + : DEFAULT_MAX_INBOUND_LINE_BYTES; + } + private getValidationErrorSignature(errors: ValidationError[]): string { if (errors.length === 0) { return 'unknown'; diff --git a/src/models/StratumV1ClientStatistics.spec.ts b/src/models/StratumV1ClientStatistics.spec.ts index 0b837f7..8702c0d 100644 --- a/src/models/StratumV1ClientStatistics.spec.ts +++ b/src/models/StratumV1ClientStatistics.spec.ts @@ -54,6 +54,22 @@ describe('StratumV1ClientStatistics', () => { expect(statistics.getSuggestedDifficulty(64)).toBe(2048); }); + it('does not retarget when accepted shares have a zero-duration sample window', async () => { + for (let i = 0; i < 5; i++) { + await statistics.addShares(client, 64); + } + + expect(() => statistics.getSuggestedDifficulty(64)).not.toThrow(); + expect(statistics.getSuggestedDifficulty(64)).toBeNull(); + }); + + it('returns a finite power-of-two difficulty for values above 32-bit range', () => { + const result = (statistics as any).nearestPowerOfTwo(2 ** 40 + 1); + + expect(result).toBe(2 ** 40); + expect(Number.isFinite(result)).toBe(true); + }); + it('should decrease difficulty for slow submissions', async () => { for (let i = 0; i < 5; i++) { jest.setSystemTime(new Date(Date.parse('2026-05-06T12:00:00Z') + (i * 150000))); diff --git a/src/models/StratumV1ClientStatistics.ts b/src/models/StratumV1ClientStatistics.ts index e56c363..8b1f751 100644 --- a/src/models/StratumV1ClientStatistics.ts +++ b/src/models/StratumV1ClientStatistics.ts @@ -51,10 +51,16 @@ export class StratumV1ClientStatistics { return pre; }, 0); const diffSeconds = (this.submissionCache[this.submissionCache.length - 1].time.getTime() - this.submissionCache[0].time.getTime()) / 1000; + if (!Number.isFinite(diffSeconds) || diffSeconds <= 0) { + return null; + } const difficultyPerSecond = sum / diffSeconds; const targetDifficulty = difficultyPerSecond * this.targetSubmitShareEveryNSeconds; + if (!Number.isFinite(targetDifficulty) || targetDifficulty <= 0) { + return null; + } if ((clientDifficulty * 2) < targetDifficulty || (clientDifficulty / 2) > targetDifficulty) { return this.nearestPowerOfTwo(targetDifficulty) @@ -63,27 +69,15 @@ export class StratumV1ClientStatistics { return null; } - private nearestPowerOfTwo(val): number { - if (val === 0) { + private nearestPowerOfTwo(val: number): number { + if (!Number.isFinite(val) || val <= 0) { return null; } - if (val < this.minDifficulty) { + if (val <= this.minDifficulty) { return this.minDifficulty; } - let x = val | (val >> 1); - x = x | (x >> 2); - x = x | (x >> 4); - x = x | (x >> 8); - x = x | (x >> 16); - x = x | (x >> 32); - const res = x - (x >> 1); - if (res == 0 && val * 100 < this.minDifficulty) { - return this.minDifficulty; - } - if (res == 0) { - return this.nearestPowerOfTwo(val * 100) / 100; - } - return res; + const result = 2 ** Math.floor(Math.log2(val)); + return Number.isFinite(result) ? Math.max(this.minDifficulty, result) : null; } }