Harden Stratum V1 share validation

This commit is contained in:
Ben
2026-08-05 16:34:43 -04:00
parent 4c56839c02
commit 07a2b173d0
5 changed files with 187 additions and 48 deletions
+1
View File
@@ -72,6 +72,7 @@ STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS=60000
# Disconnect slow/non-reading SV1 clients before their per-socket write queue # Disconnect slow/non-reading SV1 clients before their per-socket write queue
# can grow without bound. Includes bytes already buffered by Node.js. # can grow without bound. Includes bytes already buffered by Node.js.
STRATUM_MAX_SOCKET_BUFFER_BYTES=262144 STRATUM_MAX_SOCKET_BUFFER_BYTES=262144
STRATUM_MAX_INBOUND_LINE_BYTES=65536
# Immediately publish a consensus-valid, subsidy-only solo job before the full # Immediately publish a consensus-valid, subsidy-only solo job before the full
# transaction template. SV2 uses the same activation to select its pre-staged # transaction template. SV2 uses the same activation to select its pre-staged
# future job; the full job then follows on the active tip. # future job; the full job then follows on the active tip.
+91 -4
View File
@@ -169,6 +169,59 @@ describe('StratumV1Client', () => {
expect(socket.on).toHaveBeenCalled(); 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 () => { it('should clean up socket state only once when destroyed repeatedly', async () => {
const timer = setInterval(() => undefined, 1000); const timer = setInterval(() => undefined, 1000);
const removeListenerSpy = jest.spyOn(socket, 'removeListener'); const removeListenerSpy = jest.spyOn(socket, 'removeListener');
@@ -805,6 +858,7 @@ describe('StratumV1Client', () => {
await new Promise((r) => setTimeout(r, 100)); await new Promise((r) => setTimeout(r, 100));
expect((client as any).isDuplicateSubmission('old-tip-share')).toBe(false); expect((client as any).isDuplicateSubmission('old-tip-share')).toBe(false);
expect((client as any).rememberSubmission('old-tip-share', '1')).toBe(true);
const nextTip = { const nextTip = {
...MockRecording1.BLOCK_TEMPLATE, ...MockRecording1.BLOCK_TEMPLATE,
previousblockhash: '11'.repeat(32), previousblockhash: '11'.repeat(32),
@@ -817,7 +871,9 @@ describe('StratumV1Client', () => {
expect((client as any).isDuplicateSubmission('old-tip-share')).toBe(true); 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) => { (configService.get as jest.Mock).mockImplementation((key: string) => {
if (key === 'STRATUM_SUBMISSION_DEDUP_TTL_MS') return '1000'; if (key === 'STRATUM_SUBMISSION_DEDUP_TTL_MS') return '1000';
if (key === 'STRATUM_SUBMISSION_DEDUP_MAX_ENTRIES') return '2'; 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).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('two')).toBe(false);
expect((client as any).isDuplicateSubmission('three')).toBe(false); expect((client as any).rememberSubmission('two', '1')).toBe(true);
expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['two', 'three']); expect((client as any).rememberSubmission('three', '1')).toBe(false);
expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['one', 'two']);
jest.advanceTimersByTime(1001); 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).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 () => { 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((client as any).write).lastCalledWith(`{"id":5,"result":null,"error":[23,"Difficulty too low",""]}\n`);
expect(await clientService.connectedClientCount()).toBe(1); expect(await clientService.connectedClientCount()).toBe(1);
expect(shareAccountingService.recordAcceptedShare).not.toHaveBeenCalled(); 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 () => { it('rejects version rolling masks outside the negotiated BIP320 range', async () => {
+67 -26
View File
@@ -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_TTL_MS = 5 * 60 * 1000;
const DEFAULT_SUBMISSION_DEDUP_MAX_ENTRIES = 10_000; const DEFAULT_SUBMISSION_DEDUP_MAX_ENTRIES = 10_000;
const DEFAULT_MAX_SOCKET_BUFFER_BYTES = 256 * 1024; 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; const VERSION_ROLLING_MASK = 0x1fffe000;
export interface MiningJobBroadcastResult { export interface MiningJobBroadcastResult {
@@ -55,6 +57,11 @@ export interface MiningJobBroadcastResult {
preStaged?: boolean; preStaged?: boolean;
} }
interface SubmissionDedupEntry {
expiresAt: number;
jobId: string;
}
export class StratumV1Client { export class StratumV1Client {
private static blockedUserAgentLogState = new Map<string, { nextLogAt: number, suppressed: number }>(); private static blockedUserAgentLogState = new Map<string, { nextLogAt: number, suppressed: number }>();
private static validationErrorLogState = new Map<string, { nextLogAt: number, suppressed: number, sample: string }>(); private static validationErrorLogState = new Map<string, { nextLogAt: number, suppressed: number, sample: string }>();
@@ -90,8 +97,9 @@ export class StratumV1Client {
private lastHashRatePersistedAt = 0; private lastHashRatePersistedAt = 0;
private readonly network: bitcoinjs.Network; private readonly network: bitcoinjs.Network;
private readonly maxSocketBufferBytes: number; private readonly maxSocketBufferBytes: number;
private readonly maxInboundLineBytes: number;
private miningSubmissionHashes = new Map<string, number>(); private miningSubmissionHashes = new Map<string, SubmissionDedupEntry>();
constructor( constructor(
public readonly socket: Socket, public readonly socket: Socket,
@@ -115,6 +123,7 @@ export class StratumV1Client {
this.socket.on('data', this.socketDataHandler); this.socket.on('data', this.socketDataHandler);
this.network = this.getNetwork(); this.network = this.getNetwork();
this.maxSocketBufferBytes = this.readMaxSocketBufferBytes(); this.maxSocketBufferBytes = this.readMaxSocketBufferBytes();
this.maxInboundLineBytes = this.readMaxInboundLineBytes();
} }
@@ -150,9 +159,15 @@ export class StratumV1Client {
return; return;
} }
this.buffer += data.toString(); const lines = `${this.buffer}${data.toString()}`.split('\n');
const lines = this.buffer.split('\n'); const incompleteLine = lines.pop() || '';
this.buffer = lines.pop() || ''; // Save the last part of the data (incomplete line) to the buffer 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)) { for (const m of lines.filter(l => l.length > 0)) {
if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) { if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) {
@@ -679,6 +694,15 @@ export class StratumV1Client {
); );
return false; 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 // The optional BIP310 field contains replacement bits, not an XOR
// delta. A legacy five-field submission leaves the advertised version // delta. A legacy five-field submission leaves the advertised version
// unchanged, including any bits already set inside the rolling mask. // unchanged, including any bits already set inside the rolling mask.
@@ -757,6 +781,15 @@ export class StratumV1Client {
const creditedDifficulty = isBlockCandidate && !meetsSessionTarget const creditedDifficulty = isBlockCandidate && !meetsSessionTarget
? Math.min(this.sessionDifficulty, submissionDifficulty) ? Math.min(this.sessionDifficulty, submissionDifficulty)
: this.sessionDifficulty; : 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; let blockSubmissionResult: string = null;
if (status === 'stale') { if (status === 'stale') {
@@ -1052,33 +1085,31 @@ export class StratumV1Client {
private isDuplicateSubmission(submissionHash: string): boolean { private isDuplicateSubmission(submissionHash: string): boolean {
const now = Date.now(); const now = Date.now();
const existingExpiry = this.miningSubmissionHashes.get(submissionHash); this.pruneExpiredSubmissions(now);
if (existingExpiry != null && existingExpiry > now) { return this.miningSubmissionHashes.has(submissionHash);
return true;
}
if (existingExpiry != null) {
this.miningSubmissionHashes.delete(submissionHash);
} }
for (const [hash, expiresAt] of this.miningSubmissionHashes) { private rememberSubmission(submissionHash: string, jobId: string): boolean {
if (expiresAt <= now) { 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.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 { private getSubmissionDedupTtlMs(): number {
@@ -1111,6 +1142,16 @@ export class StratumV1Client {
: DEFAULT_MAX_SOCKET_BUFFER_BYTES; : DEFAULT_MAX_SOCKET_BUFFER_BYTES;
} }
private readMaxInboundLineBytes(): number {
const configured = Number(
this.configService.get<string>('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 { private getValidationErrorSignature(errors: ValidationError[]): string {
if (errors.length === 0) { if (errors.length === 0) {
return 'unknown'; return 'unknown';
@@ -54,6 +54,22 @@ describe('StratumV1ClientStatistics', () => {
expect(statistics.getSuggestedDifficulty(64)).toBe(2048); 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 () => { it('should decrease difficulty for slow submissions', async () => {
for (let i = 0; i < 5; i++) { for (let i = 0; i < 5; i++) {
jest.setSystemTime(new Date(Date.parse('2026-05-06T12:00:00Z') + (i * 150000))); jest.setSystemTime(new Date(Date.parse('2026-05-06T12:00:00Z') + (i * 150000)));
+11 -17
View File
@@ -51,10 +51,16 @@ export class StratumV1ClientStatistics {
return pre; return pre;
}, 0); }, 0);
const diffSeconds = (this.submissionCache[this.submissionCache.length - 1].time.getTime() - this.submissionCache[0].time.getTime()) / 1000; 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 difficultyPerSecond = sum / diffSeconds;
const targetDifficulty = difficultyPerSecond * this.targetSubmitShareEveryNSeconds; const targetDifficulty = difficultyPerSecond * this.targetSubmitShareEveryNSeconds;
if (!Number.isFinite(targetDifficulty) || targetDifficulty <= 0) {
return null;
}
if ((clientDifficulty * 2) < targetDifficulty || (clientDifficulty / 2) > targetDifficulty) { if ((clientDifficulty * 2) < targetDifficulty || (clientDifficulty / 2) > targetDifficulty) {
return this.nearestPowerOfTwo(targetDifficulty) return this.nearestPowerOfTwo(targetDifficulty)
@@ -63,27 +69,15 @@ export class StratumV1ClientStatistics {
return null; return null;
} }
private nearestPowerOfTwo(val): number { private nearestPowerOfTwo(val: number): number {
if (val === 0) { if (!Number.isFinite(val) || val <= 0) {
return null; return null;
} }
if (val < this.minDifficulty) { if (val <= this.minDifficulty) {
return this.minDifficulty; return this.minDifficulty;
} }
let x = val | (val >> 1); const result = 2 ** Math.floor(Math.log2(val));
x = x | (x >> 2); return Number.isFinite(result) ? Math.max(this.minDifficulty, result) : null;
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;
} }
} }