From 216fa74d99bfad5ffc16ebc91845104c4fd32d67 Mon Sep 17 00:00:00 2001 From: Ben Date: Wed, 6 May 2026 19:05:10 -0400 Subject: [PATCH] misc improvements from blitzpool --- .../address-settings.service.ts | 21 ++-- src/ORM/utils/DateTimeTransformer.spec.ts | 38 ++++++ src/ORM/utils/DateTimeTransformer.ts | 115 ++++++++++++++++-- src/models/MiningJob.spec.ts | 28 +++++ src/models/MiningJob.ts | 38 +++++- src/models/StratumV1Client.spec.ts | 67 +++++++++- src/models/StratumV1Client.ts | 85 ++++++++++--- src/models/StratumV1ClientStatistics.spec.ts | 6 +- src/models/StratumV1ClientStatistics.ts | 44 ++++--- .../MiningSubmitMessage.spec.ts | 23 ++++ .../stratum-messages/MiningSubmitMessage.ts | 10 +- src/services/bitcoin-rpc.service.ts | 12 +- src/services/stratum-v1-jobs.service.spec.ts | 42 +++++-- src/services/stratum-v1-jobs.service.ts | 98 +++++++++++---- 14 files changed, 526 insertions(+), 101 deletions(-) create mode 100644 src/ORM/utils/DateTimeTransformer.spec.ts diff --git a/src/ORM/address-settings/address-settings.service.ts b/src/ORM/address-settings/address-settings.service.ts index 884050b..5f2a558 100644 --- a/src/ORM/address-settings/address-settings.service.ts +++ b/src/ORM/address-settings/address-settings.service.ts @@ -15,14 +15,17 @@ export class AddressSettingsService { } public async getSettings(address: string, createIfNotFound: boolean) { - const settings = await this.addressSettingsRepository.findOne({ where: { address } }); - if (createIfNotFound == true && settings == null) { - // It's possible to have a race condition here so if we get a PK violation, fetch it - try { - return await this.createNew(address); - } catch (e) { - return await this.addressSettingsRepository.findOne({ where: { address } }); - } + let settings = await this.addressSettingsRepository.findOne({ where: { address } }); + if (createIfNotFound === true && settings == null) { + await this.addressSettingsRepository + .createQueryBuilder() + .insert() + .into(AddressSettingsEntity) + .values({ address }) + .orIgnore() + .execute(); + + settings = await this.addressSettingsRepository.findOne({ where: { address } }); } return settings; } @@ -59,4 +62,4 @@ export class AddressSettingsService { .limit(10) .execute(); } -} \ No newline at end of file +} diff --git a/src/ORM/utils/DateTimeTransformer.spec.ts b/src/ORM/utils/DateTimeTransformer.spec.ts new file mode 100644 index 0000000..c23651e --- /dev/null +++ b/src/ORM/utils/DateTimeTransformer.spec.ts @@ -0,0 +1,38 @@ +import { DateTimeTransformer } from './DateTimeTransformer'; + +describe('DateTimeTransformer', () => { + const transformer = new DateTimeTransformer(); + + it('should pass Date values through as Date instances', () => { + const date = new Date('2026-05-06T12:00:00.000Z'); + + expect(transformer.to(date)).toBe(date); + expect(transformer.from(date)).toBe(date); + }); + + it('should parse ISO strings from the database', () => { + const date = transformer.from('2026-05-06T12:00:00.000Z'); + + expect(date).toBeInstanceOf(Date); + expect((date as Date).toISOString()).toBe('2026-05-06T12:00:00.000Z'); + }); + + it('should parse legacy locale strings', () => { + const usDate = transformer.from('5/6/2026, 1:02:03 PM') as Date; + const europeanDate = transformer.from('6.5.2026, 13:02:03') as Date; + + expect(usDate.getFullYear()).toBe(2026); + expect(usDate.getMonth()).toBe(4); + expect(usDate.getDate()).toBe(6); + expect(usDate.getHours()).toBe(13); + expect(europeanDate.getFullYear()).toBe(2026); + expect(europeanDate.getMonth()).toBe(4); + expect(europeanDate.getDate()).toBe(6); + expect(europeanDate.getHours()).toBe(13); + }); + + it('should preserve null and undefined values', () => { + expect(transformer.to(null)).toBeNull(); + expect(transformer.from(undefined)).toBeUndefined(); + }); +}); diff --git a/src/ORM/utils/DateTimeTransformer.ts b/src/ORM/utils/DateTimeTransformer.ts index fa2c07e..0d66ff8 100644 --- a/src/ORM/utils/DateTimeTransformer.ts +++ b/src/ORM/utils/DateTimeTransformer.ts @@ -1,14 +1,113 @@ import { ValueTransformer } from 'typeorm'; -export class DateTimeTransformer implements ValueTransformer { - to(value: Date): any { - // Convert the local time to UTC before saving to the database - const utcTime = value?.toLocaleString(); - return utcTime; +const EUROPEAN_LOCALE_REGEX = + /^(\d{1,2})\.(\d{1,2})\.(\d{4}),\s*(\d{1,2}):(\d{2})(?::(\d{2}))?(?:\s*(AM|PM))?$/i; +const US_LOCALE_REGEX = + /^(\d{1,2})\/(\d{1,2})\/(\d{4}),\s*(\d{1,2}):(\d{2})(?::(\d{2}))?\s*(AM|PM)$/i; + +function parseLegacyLocaleString(value: string): Date | null { + const trimmed = value.trim(); + + const matchEuropean = trimmed.match(EUROPEAN_LOCALE_REGEX); + if (matchEuropean) { + const [, day, month, year, hours, minutes, seconds = '0', meridiem] = matchEuropean; + return buildDateFromParts({ + year, + month, + day, + hours, + minutes, + seconds, + meridiem, + }); } - from(value: any): Date { - // Convert the UTC time from the database to the local time zone + const matchUs = trimmed.match(US_LOCALE_REGEX); + if (matchUs) { + const [, month, day, year, hours, minutes, seconds = '0', meridiem] = matchUs; + return buildDateFromParts({ + year, + month, + day, + hours, + minutes, + seconds, + meridiem, + }); + } + + return null; +} + +function buildDateFromParts({ + year, + month, + day, + hours, + minutes, + seconds, + meridiem, +}: { + year: string; + month: string; + day: string; + hours: string; + minutes: string; + seconds: string; + meridiem?: string; +}): Date { + let hourValue = Number(hours); + if (meridiem) { + const normalized = meridiem.toUpperCase(); + if (normalized === 'PM' && hourValue < 12) { + hourValue += 12; + } else if (normalized === 'AM' && hourValue === 12) { + hourValue = 0; + } + } + + return new Date( + Number(year), + Number(month) - 1, + Number(day), + hourValue, + Number(minutes), + Number(seconds), + ); +} + +function normalizeDate(value: Date | string): Date { + if (value instanceof Date) { return value; } -} \ No newline at end of file + + const legacyDate = parseLegacyLocaleString(value); + if (legacyDate) { + return legacyDate; + } + + const date = new Date(value); + if (!Number.isNaN(date.getTime())) { + return date; + } + + throw new Error(`Invalid date value received: ${value}`); +} + +export class DateTimeTransformer implements ValueTransformer { + to(value: Date | string | null | undefined): Date | null | undefined { + if (value === null || value === undefined) { + return value === undefined ? undefined : null; + } + + return normalizeDate(value); + } + + from(value: Date | string | null | undefined): Date | null | undefined { + if (value === null || value === undefined) { + return value === undefined ? undefined : null; + } + + return normalizeDate(value); + } +} diff --git a/src/models/MiningJob.spec.ts b/src/models/MiningJob.spec.ts index f7d392f..2644014 100644 --- a/src/models/MiningJob.spec.ts +++ b/src/models/MiningJob.spec.ts @@ -82,4 +82,32 @@ describe('MiningJob', () => { expect(updatedBlock.version).toBe(jobTemplate.block.version); }); + + it('should expose coinbase helpers for downstream validation', () => { + const notify = JSON.parse(job.response(jobTemplate)); + + expect(job.getCoinbaseTxHex()).toContain(notify.params[2]); + expect(job.getCoinbasePrefixBuffer().toString('hex')).toBe(notify.params[2]); + expect(job.getCoinbaseSuffixBuffer().toString('hex')).toBe(notify.params[3]); + const clone = job.cloneCoinbaseTransaction(); + const nonWitnessCoinbase = bitcoinjs.Transaction.fromHex(job.getCoinbaseTxHex()); + expect(clone.ins[0].script.toString('hex')).toBe(nonWitnessCoinbase.ins[0].script.toString('hex')); + expect(clone.outs).toHaveLength(nonWitnessCoinbase.outs.length); + expect(clone.ins[0].witness[0]).toHaveLength(32); + }); + + it('should omit oversized pool identifiers from the coinbase script', () => { + jest.spyOn(console, 'warn').mockImplementation(() => undefined); + const oversizedIdentifier = 'x'.repeat(120); + const oversizedJob = new MiningJob( + bitcoinjs.networks.testnet, + '2', + [{ address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', percent: 100 }], + jobTemplate, + oversizedIdentifier + ); + const coinbase = bitcoinjs.Transaction.fromHex(oversizedJob.getCoinbaseTxHex()); + + expect(coinbase.ins[0].script.toString()).not.toContain(oversizedIdentifier); + }); }); diff --git a/src/models/MiningJob.ts b/src/models/MiningJob.ts index 2def743..1fb458b 100644 --- a/src/models/MiningJob.ts +++ b/src/models/MiningJob.ts @@ -6,6 +6,8 @@ import { eResponseMethod } from './enums/eResponseMethod'; import { IMiningNotify } from './stratum-messages/IMiningNotify'; import { TOTAL_EXTRANONCE_SIZE_BYTES } from './stratum.constants'; +const MAX_BLOCK_WEIGHT = 4000000; +const MAX_SCRIPT_SIZE = 100; interface AddressObject { address: string; @@ -20,12 +22,14 @@ export class MiningJob { public jobTemplateId: string; public networkDifficulty: number; public creation: number; + public retiredAt?: number; constructor( private network: bitcoinjs.networks.Network, public jobId: string, payoutInformation: AddressObject[], - jobTemplate: IJobTemplate + jobTemplate: IJobTemplate, + poolIdentifier: string = process.env.POOL_IDENTIFIER || 'Public-Pool' ) { this.creation = new Date().getTime(); @@ -41,7 +45,7 @@ export class MiningJob { // 32-byte - Commitment hash: Double-SHA256(witness root hash|witness reserved value) // 39th byte onwards: Optional data with no consensus meaning - const extra = Buffer.from('Public-Pool'); + const extra = Buffer.from(poolIdentifier); // Encode the block height // https://github.com/bitcoin/bips/blob/master/bip-0034.mediawiki @@ -54,10 +58,21 @@ export class MiningJob { const padding = Buffer.alloc(TOTAL_EXTRANONCE_SIZE_BYTES + (3 - blockHeightEncoded.length), 0) // build the script - this.coinbaseTransaction.ins[0].script = Buffer.concat([blockHeightLengthByte, blockHeightEncoded, extra, padding]) + let script = Buffer.concat([blockHeightLengthByte, blockHeightEncoded, extra, padding]); + if (script.length > MAX_SCRIPT_SIZE) { + console.warn('Pool identifier is too long, removing the pool identifier'); + script = Buffer.concat([blockHeightLengthByte, blockHeightEncoded, padding]); + } + + this.coinbaseTransaction.ins[0].script = script; this.coinbaseTransaction.addOutput(bitcoinjs.script.compile([bitcoinjs.opcodes.OP_RETURN, Buffer.concat([segwitMagicBits, jobTemplate.block.witnessCommit])]), 0); + if ((this.coinbaseTransaction.weight() + jobTemplate.block.weight()) > MAX_BLOCK_WEIGHT) { + console.warn('Block weight exceeds the maximum allowed weight, removing the pool identifier'); + this.coinbaseTransaction.ins[0].script = Buffer.concat([blockHeightLengthByte, blockHeightEncoded, padding]); + } + // get the non-witness coinbase tx //@ts-ignore const serializedCoinbaseTx = this.coinbaseTransaction.__toBuffer().toString('hex'); @@ -72,6 +87,23 @@ export class MiningJob { } + public getCoinbaseTxHex(): string { + //@ts-ignore + return this.coinbaseTransaction.__toBuffer().toString('hex'); + } + + public getCoinbasePrefixBuffer(): Buffer { + return Buffer.from(this.coinbasePart1, 'hex'); + } + + public getCoinbaseSuffixBuffer(): Buffer { + return Buffer.from(this.coinbasePart2, 'hex'); + } + + public cloneCoinbaseTransaction(): bitcoinjs.Transaction { + return bitcoinjs.Transaction.fromBuffer(this.coinbaseTransaction.toBuffer()); + } + public copyAndUpdateBlock(jobTemplate: IJobTemplate, versionMask: number, nonce: number, extraNonce: string, extraNonce2: string, timestamp: number): bitcoinjs.Block { const testBlock = Object.assign(new bitcoinjs.Block(), jobTemplate.block); diff --git a/src/models/StratumV1Client.spec.ts b/src/models/StratumV1Client.spec.ts index 94d6428..aa062fd 100644 --- a/src/models/StratumV1Client.spec.ts +++ b/src/models/StratumV1Client.spec.ts @@ -20,7 +20,7 @@ import { BitcoinRpcService as MockBitcoinRpcService } from '../services/bitcoin- import { NotificationService } from '../services/notification.service'; import { StratumV1JobsService } from '../services/stratum-v1-jobs.service'; import { IBlockTemplate } from './bitcoin-rpc/IBlockTemplate'; -import { StratumV1Client } from './StratumV1Client'; +import { effectiveJobDifficulty, StratumV1Client } from './StratumV1Client'; @@ -54,6 +54,7 @@ describe('StratumV1Client', () => { const emitMessage = (message: string) => socketEmitter(Buffer.from(`${message}\n`)); let consoleLogSpy: jest.SpyInstance; let consoleErrorSpy: jest.SpyInstance; + let consoleWarnSpy: jest.SpyInstance; let newBlockEmitter: BehaviorSubject = new BehaviorSubject(MockRecording1.BLOCK_TEMPLATE); @@ -103,6 +104,7 @@ describe('StratumV1Client', () => { jest.setSystemTime(new Date(parseInt(MockRecording1.TIME, 16) * 1000)); consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined); consoleErrorSpy = jest.spyOn(console, 'error').mockImplementation(() => undefined); + consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined); clientService = moduleRef.get(ClientService); @@ -168,6 +170,7 @@ describe('StratumV1Client', () => { client.destroy(); consoleLogSpy.mockRestore(); consoleErrorSpy.mockRestore(); + consoleWarnSpy.mockRestore(); jest.useRealTimers(); }) @@ -182,6 +185,27 @@ describe('StratumV1Client', () => { expect(socket.on).toHaveBeenCalled(); }); + it('should process batched socket messages sequentially', async () => { + const calls: string[] = []; + jest.spyOn(client as any, 'handleMessage').mockImplementation(async (message: string) => { + calls.push(`start:${message}`); + await Promise.resolve(); + calls.push(`end:${message}`); + }); + + socketEmitter(Buffer.from('first\nsecond\n')); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + + expect(calls).toEqual([ + 'start:first', + 'end:first', + 'start:second', + 'end:second', + ]); + }); + it('should respond to mining.subscribe', async () => { jest.spyOn(socket, 'write').mockImplementation((data) => true); @@ -371,12 +395,14 @@ describe('StratumV1Client', () => { it('should submit and persist found blocks', async () => { jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); + (bitcoinRpcService.SUBMIT_BLOCK as jest.Mock).mockResolvedValue('SUCCESS!'); jest.spyOn(client as any, 'calculateDifficulty').mockReturnValue({ submissionDifficulty: Number.MAX_SAFE_INTEGER, submissionHash: 'block-share' }); const addressSettings = moduleRef.get(AddressSettingsService); - jest.spyOn(addressSettings, 'resetBestDifficultyAndShares').mockResolvedValue(undefined); + const resetSpy = jest.spyOn(addressSettings, 'resetBestDifficultyAndShares').mockResolvedValue(undefined); + resetSpy.mockClear(); emitMessage(MockRecording1.MINING_SUBSCRIBE); emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`); @@ -400,6 +426,43 @@ describe('StratumV1Client', () => { expect((client as any).write).lastCalledWith(`{"id":5,"error":null,"result":true}\n`); }); + it('should not persist or notify rejected block submissions', async () => { + jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); + (bitcoinRpcService.SUBMIT_BLOCK as jest.Mock).mockResolvedValue('bad-prevblk'); + jest.spyOn(client as any, 'calculateDifficulty').mockReturnValue({ + submissionDifficulty: Number.MAX_SAFE_INTEGER, + submissionHash: 'block-share' + }); + const addressSettings = moduleRef.get(AddressSettingsService); + const resetSpy = jest.spyOn(addressSettings, 'resetBestDifficultyAndShares').mockResolvedValue(undefined); + resetSpy.mockClear(); + + emitMessage(MockRecording1.MINING_SUBSCRIBE); + emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`); + emitMessage(MockRecording1.MINING_AUTHORIZE); + await new Promise((r) => setTimeout(r, 100)); + + emitMessage(MockRecording1.MINING_SUBMIT); + jest.useRealTimers(); + await new Promise((r) => setTimeout(r, 1000)); + + expect(bitcoinRpcService.SUBMIT_BLOCK).toHaveBeenCalledWith(expect.any(String)); + expect(blocksService.save).not.toHaveBeenCalled(); + expect(notificationService.notifySubscribersBlockFound).not.toHaveBeenCalled(); + expect(resetSpy).not.toHaveBeenCalled(); + expect((client as any).write).lastCalledWith(`{"id":5,"error":null,"result":true}\n`); + }); + }); + +describe('effectiveJobDifficulty', () => { + it('should use the old difficulty for jobs issued before a difficulty change', () => { + expect(effectiveJobDifficulty(9, 128, 32, 10)).toBe(32); + }); + + it('should use the current difficulty for jobs issued after a difficulty change', () => { + expect(effectiveJobDifficulty(10, 128, 32, 10)).toBe(128); + }); +}); diff --git a/src/models/StratumV1Client.ts b/src/models/StratumV1Client.ts index 5e91468..afdd313 100644 --- a/src/models/StratumV1Client.ts +++ b/src/models/StratumV1Client.ts @@ -29,6 +29,21 @@ import { EXTRANONCE1_SIZE_BYTES } from './stratum.constants'; import { SuggestDifficulty } from './stratum-messages/SuggestDifficultyMessage'; import { StratumV1ClientStatistics } from './StratumV1ClientStatistics'; +export function effectiveJobDifficulty( + jobIdInt: number, + currentDiff: number, + oldDiff: number, + diffChangeJobId: number | null, +): number { + if (diffChangeJobId == null || !Number.isFinite(jobIdInt)) { + return currentDiff; + } + if (jobIdInt < diffChangeJobId) { + return Math.min(currentDiff, oldDiff); + } + return currentDiff; +} + export class StratumV1Client { @@ -43,6 +58,8 @@ export class StratumV1Client { private stratumInitialized = false; private usedSuggestedDifficulty = false; private sessionDifficulty: number = 100000; + private oldSessionDifficulty: number = 100000; + private diffChangeJobId: number | null = null; private clientEntity: ClientEntity; private creatingEntity: Promise; @@ -73,16 +90,16 @@ export class StratumV1Client { let lines = this.buffer.split('\n'); this.buffer = lines.pop() || ''; // Save the last part of the data (incomplete line) to the buffer - lines - .filter(m => m.length > 0) - .forEach(async (m) => { + (async () => { + for (const m of lines.filter(l => l.length > 0)) { try { await this.handleMessage(m); } catch (e) { await this.socket.end(); console.error(e); } - }); + } + })(); }); @@ -224,6 +241,7 @@ export class StratumV1Client { this.clientAuthorization = authorizationMessage; if (this.clientSuggestedDifficulty == null && this.clientAuthorization.startingDiff != null && this.clientAuthorization.startingDiff > this.sessionDifficulty) { this.sessionDifficulty = this.clientAuthorization.startingDiff; + this.oldSessionDifficulty = this.sessionDifficulty; } const success = await this.write(JSON.stringify(this.clientAuthorization.response()) + '\n'); if (!success) { @@ -265,6 +283,7 @@ export class StratumV1Client { this.clientSuggestedDifficulty = suggestDifficultyMessage; this.sessionDifficulty = suggestDifficultyMessage.suggestedDifficulty; + this.oldSessionDifficulty = this.sessionDifficulty; const success = await this.write(JSON.stringify(this.clientSuggestedDifficulty.response(this.sessionDifficulty)) + '\n'); if (!success) { return; @@ -361,6 +380,7 @@ export class StratumV1Client { switch (this.clientSubscription.userAgent) { case 'cpuminer': { this.sessionDifficulty = 0.1; + this.oldSessionDifficulty = this.sessionDifficulty; } } @@ -443,7 +463,8 @@ export class StratumV1Client { network, this.stratumV1JobsService.getNextId(), payoutInformation, - jobTemplate + jobTemplate, + this.configService.get('POOL_IDENTIFIER') || 'Public-Pool' ); this.stratumV1JobsService.addJob(job); @@ -531,6 +552,19 @@ export class StratumV1Client { } return false; } + + const classification = this.stratumV1JobsService.classifyJobForShare(job); + if (classification === 'stale-rejected') { + const err = new StratumErrorMessage( + submission.id, + eStratumErrorCode.JobNotFound, + 'stale').response(); + const success = await this.write(err); + if (!success) { + return false; + } + return false; + } const updatedJobBlock = job.copyAndUpdateBlock( @@ -545,30 +579,37 @@ export class StratumV1Client { const { submissionDifficulty } = this.calculateDifficulty(header); //console.log(`DIFF: ${submissionDifficulty} of ${this.sessionDifficulty} from ${this.clientAuthorization.worker + '.' + this.extraNonceAndSessionId}`); + const effectiveDiff = effectiveJobDifficulty( + parseInt(submission.jobId, 16), + this.sessionDifficulty, + this.oldSessionDifficulty, + this.diffChangeJobId, + ); - if (submissionDifficulty >= this.sessionDifficulty) { + if (submissionDifficulty >= effectiveDiff) { if (submissionDifficulty >= jobTemplate.blockData.networkDifficulty) { console.log('!!! BLOCK FOUND !!!'); const blockHex = updatedJobBlock.toHex(false); const result = await this.bitcoinRpcService.SUBMIT_BLOCK(blockHex); - await this.blocksService.save({ - height: jobTemplate.blockData.height, - minerAddress: this.clientAuthorization.address, - worker: this.clientAuthorization.worker, - sessionId: this.extraNonceAndSessionId, - blockData: blockHex - }); + if (result === 'SUCCESS!') { + await this.blocksService.save({ + height: jobTemplate.blockData.height, + minerAddress: this.clientAuthorization.address, + worker: this.clientAuthorization.worker, + sessionId: this.extraNonceAndSessionId, + blockData: blockHex + }); - await this.notificationService.notifySubscribersBlockFound(this.clientAuthorization.address, jobTemplate.blockData.height, updatedJobBlock, result); - //success - if (result == null) { + await this.notificationService.notifySubscribersBlockFound(this.clientAuthorization.address, jobTemplate.blockData.height, updatedJobBlock, result); await this.addressSettingsService.resetBestDifficultyAndShares(); + } else { + console.warn(`[Block submit rejected at height ${jobTemplate.blockData.height}]: ${result}`); } } try { - await this.statistics.addShares(this.clientEntity, this.sessionDifficulty); + await this.statistics.addShares(this.clientEntity, effectiveDiff); const now = new Date(); // only update every minute //if (this.clientEntity.updatedAt == null || now.getTime() - this.clientEntity.updatedAt.getTime() > 1000 * 60) { @@ -617,6 +658,11 @@ export class StratumV1Client { if (targetDiff != this.sessionDifficulty) { //console.log(`Adjusting ${this.extraNonceAndSessionId} difficulty from ${this.sessionDifficulty} to ${targetDiff}`); + if (!Number.isFinite(targetDiff)) { + return; + } + this.oldSessionDifficulty = this.sessionDifficulty; + this.diffChangeJobId = parseInt(this.stratumV1JobsService.getNextId(), 16); this.sessionDifficulty = targetDiff; const data = JSON.stringify({ @@ -626,7 +672,10 @@ export class StratumV1Client { }) + '\n'; - await this.socket.write(data); + const success = await this.write(data); + if (!success) { + return; + } const jobTemplate = await firstValueFrom(this.stratumV1JobsService.newMiningJob$); // we need to clear the jobs so that the difficulty set takes effect. Otherwise the different miner implementations can cause issues diff --git a/src/models/StratumV1ClientStatistics.spec.ts b/src/models/StratumV1ClientStatistics.spec.ts index 33cf0ec..2e7633e 100644 --- a/src/models/StratumV1ClientStatistics.spec.ts +++ b/src/models/StratumV1ClientStatistics.spec.ts @@ -81,10 +81,10 @@ describe('StratumV1ClientStatistics', () => { expect(statistics.getSuggestedDifficulty(64)).toBeNull(); }); - it('should lower difficulty when a miner has not submitted shares for several minutes', () => { - jest.setSystemTime(new Date('2026-05-06T12:06:00Z')); + it('should lower difficulty when a miner has not submitted shares for more than a minute', () => { + jest.setSystemTime(new Date('2026-05-06T12:01:01Z')); - expect(statistics.getSuggestedDifficulty(64)).toBe(8); + expect(statistics.getSuggestedDifficulty(64)).toBe(12); }); it('should increase difficulty for rapid submissions', async () => { diff --git a/src/models/StratumV1ClientStatistics.ts b/src/models/StratumV1ClientStatistics.ts index aebe42d..e257e5a 100644 --- a/src/models/StratumV1ClientStatistics.ts +++ b/src/models/StratumV1ClientStatistics.ts @@ -104,8 +104,8 @@ export class StratumV1ClientStatistics { // miner hasn't submitted shares in one minute if (this.submissionCache.length < 5) { - if ((new Date().getTime() - this.submissionCacheStart.getTime()) / 5000 > 60) { - return this.nearestPowerOfTwo(clientDifficulty / 6); + if ((new Date().getTime() - this.submissionCacheStart.getTime()) / 1000 > 60) { + return this.nearestDifficultyStep(clientDifficulty / 6); } else { return null; } @@ -116,39 +116,45 @@ export class StratumV1ClientStatistics { return pre; }, 0); const diffSeconds = (this.submissionCache[this.submissionCache.length - 1].time.getTime() - this.submissionCache[0].time.getTime()) / 1000; + if (diffSeconds <= 0) { + return null; + } const difficultyPerSecond = sum / diffSeconds; const targetDifficulty = difficultyPerSecond * this.targetSubmitShareEveryNSeconds; + if (!Number.isFinite(difficultyPerSecond) || !Number.isFinite(targetDifficulty)) { + return null; + } if ((clientDifficulty * 2) < targetDifficulty || (clientDifficulty / 2) > targetDifficulty) { - return this.nearestPowerOfTwo(targetDifficulty) + return this.nearestDifficultyStep(targetDifficulty) } return null; } - private nearestPowerOfTwo(val): number { + private nearestDifficultyStep(val: number): number { if (val === 0) { return null; } if (val < MIN_DIFF) { return MIN_DIFF; } - 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 < MIN_DIFF) { - return MIN_DIFF; - } - if (res == 0) { - return this.nearestPowerOfTwo(val * 100) / 100; - } - return res; + + const exponent = Math.floor(Math.log2(val)); + const lower = 2 ** exponent; + const middle = lower + lower / 2; + const upper = lower * 2; + + const distances = [ + { value: lower, diff: Math.abs(val - lower) }, + { value: middle, diff: Math.abs(val - middle) }, + { value: upper, diff: Math.abs(val - upper) }, + ]; + + distances.sort((a, b) => a.diff - b.diff); + return distances[0].value; } -} \ No newline at end of file +} diff --git a/src/models/stratum-messages/MiningSubmitMessage.spec.ts b/src/models/stratum-messages/MiningSubmitMessage.spec.ts index 01fa3e1..d196314 100644 --- a/src/models/stratum-messages/MiningSubmitMessage.spec.ts +++ b/src/models/stratum-messages/MiningSubmitMessage.spec.ts @@ -58,6 +58,29 @@ describe('MiningSubmitMessage', () => { expect(errors.some(error => error.property === 'extraNonce2')).toBe(true); }); + + it('should hash submit fields without concatenation ambiguity', () => { + const first = plainToInstance( + MiningSubmitMessage, + { + id: 5, + method: 'mining.submit', + params: ['worker', 'e', '9902000000000000', 'd', 'bc', 'a'] + }, + ); + const second = plainToInstance( + MiningSubmitMessage, + { + id: 5, + method: 'mining.submit', + params: ['worker', 'e', '9902000000000000', 'd', 'c', 'ab'] + }, + ); + + expect(first.versionMask + first.nonce + first.extraNonce2 + first.ntime + first.jobId) + .toBe(second.versionMask + second.nonce + second.extraNonce2 + second.ntime + second.jobId); + expect(first.hash()).not.toBe(second.hash()); + }); }); diff --git a/src/models/stratum-messages/MiningSubmitMessage.ts b/src/models/stratum-messages/MiningSubmitMessage.ts index 8fbd121..e76d501 100644 --- a/src/models/stratum-messages/MiningSubmitMessage.ts +++ b/src/models/stratum-messages/MiningSubmitMessage.ts @@ -68,8 +68,14 @@ export class MiningSubmitMessage extends StratumBaseMessage { } public hash(): string{ - const buffer = Buffer.from(this.versionMask + this.nonce + this.extraNonce2 + this.ntime + this.jobId); - return bitcoinjs.crypto.hash256(buffer).toString('base64'); + const canonical = JSON.stringify({ + versionMask: this.versionMask ?? '', + nonce: this.nonce, + extraNonce2: this.extraNonce2, + ntime: this.ntime, + jobId: this.jobId, + }); + return bitcoinjs.crypto.hash256(Buffer.from(canonical)).toString('base64'); } diff --git a/src/services/bitcoin-rpc.service.ts b/src/services/bitcoin-rpc.service.ts index d58e4c3..a1c9570 100644 --- a/src/services/bitcoin-rpc.service.ts +++ b/src/services/bitcoin-rpc.service.ts @@ -8,6 +8,7 @@ import * as zmq from 'zeromq'; import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate'; import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo'; import * as PGPubsub from 'pg-pubsub'; +import * as fs from 'node:fs'; @Injectable() export class BitcoinRpcService implements OnModuleInit { @@ -34,10 +35,17 @@ export class BitcoinRpcService implements OnModuleInit { const url = this.configService.get('BITCOIN_RPC_URL'); - const user = this.configService.get('BITCOIN_RPC_USER'); - const pass = this.configService.get('BITCOIN_RPC_PASSWORD'); + let user = this.configService.get('BITCOIN_RPC_USER'); + let pass = this.configService.get('BITCOIN_RPC_PASSWORD'); const port = parseInt(this.configService.get('BITCOIN_RPC_PORT')); const timeout = parseInt(this.configService.get('BITCOIN_RPC_TIMEOUT')); + const cookiefile = this.configService.get('BITCOIN_RPC_COOKIEFILE'); + + if (cookiefile != null && cookiefile !== '') { + const [cookieUser, cookiePass] = fs.readFileSync(cookiefile).toString().trim().split(':'); + user = cookieUser; + pass = cookiePass; + } this.client = new RPCClient({ url, port, timeout, user, pass }); diff --git a/src/services/stratum-v1-jobs.service.spec.ts b/src/services/stratum-v1-jobs.service.spec.ts index d35b935..de5cbfc 100644 --- a/src/services/stratum-v1-jobs.service.spec.ts +++ b/src/services/stratum-v1-jobs.service.spec.ts @@ -47,8 +47,8 @@ describe('StratumV1JobsService', () => { expect(service.getJobTemplateById('1')).toBe(jobTemplate); }); - it('should clear stale jobs when the block height changes', async () => { - await firstValueFrom(service.newMiningJob$); + it('should retire jobs when the block height changes', async () => { + const firstTemplate = await firstValueFrom(service.newMiningJob$); service.addJob({ jobId: 'old-job', creation: Date.now() } as any); bitcoinRpcService.miningInfo.blocks = MockRecording1.BLOCK_TEMPLATE.height + 1; @@ -57,17 +57,24 @@ describe('StratumV1JobsService', () => { const jobTemplate = await nextTemplate; expect(jobTemplate.blockData.clearJobs).toBe(true); - expect(service.getJobById('old-job')).toBeUndefined(); - expect(Object.keys(service.blocks)).toEqual([jobTemplate.blockData.id]); + expect(service.getJobById('old-job')).toEqual(expect.objectContaining({ + jobId: 'old-job', + retiredAt: Date.now() + })); + expect(service.getJobTemplateById(firstTemplate.blockData.id).blockData.retiredAt).toBe(Date.now()); + expect(service.getJobTemplateById(jobTemplate.blockData.id)).toBe(jobTemplate); }); - it('should expire old jobs and templates when the block height is unchanged', async () => { + it('should age retired jobs and templates after the retention window', async () => { await firstValueFrom(service.newMiningJob$); - const oldCreation = Date.now() - (1000 * 60 * 6); - service.jobs['old-job'] = { jobId: 'old-job', creation: oldCreation } as any; - service.blocks['old-template'] = { - blockData: { creation: oldCreation } - } as any; + const oldCreation = Date.now() - (1000 * 60 * 11); + const retiredAt = Date.now() - (1000 * 60 * 11); + for (let i = 0; i < 5; i++) { + service.jobs[`old-job-${i}`] = { jobId: `old-job-${i}`, creation: oldCreation - i, retiredAt } as any; + service.blocks[`old-template-${i}`] = { + blockData: { creation: oldCreation - i, retiredAt } + } as any; + } bitcoinRpcService.miningInfo.blocks = MockRecording1.BLOCK_TEMPLATE.height; const nextTemplate = firstValueFrom(service.newMiningJob$.pipe(skip(1))); @@ -75,11 +82,22 @@ describe('StratumV1JobsService', () => { const jobTemplate = await nextTemplate; expect(jobTemplate.blockData.clearJobs).toBe(false); - expect(service.getJobById('old-job')).toBeUndefined(); - expect(service.getJobTemplateById('old-template')).toBeUndefined(); + expect(service.getJobById('old-job-4')).toBeUndefined(); + expect(service.getJobTemplateById('old-template-4')).toBeUndefined(); expect(service.getJobTemplateById(jobTemplate.blockData.id)).toBe(jobTemplate); }); + it('should classify retired jobs inside and outside the stale grace window', () => { + const job = { jobId: '1', creation: Date.now(), retiredAt: Date.now() - 1000 } as any; + + expect(service.classifyJobForShare(job, Date.now())).toBe('stale-creditable'); + + job.retiredAt = Date.now() - 6000; + + expect(service.classifyJobForShare(job, Date.now())).toBe('stale-rejected'); + expect(service.classifyJobForShare({ jobId: '2', creation: Date.now() } as any, Date.now())).toBe('active'); + }); + it('should increment job ids when jobs are added', () => { expect(service.getNextId()).toBe('1'); diff --git a/src/services/stratum-v1-jobs.service.ts b/src/services/stratum-v1-jobs.service.ts index 37e1246..40e5c29 100644 --- a/src/services/stratum-v1-jobs.service.ts +++ b/src/services/stratum-v1-jobs.service.ts @@ -18,9 +18,14 @@ export interface IJobTemplate { networkDifficulty: number; height: number; clearJobs: boolean; + retiredAt?: number; }; } +const STALE_GRACE_MS = parseInt(process.env.STRATUM_STALE_GRACE_MS) || 5000; +const MIN_RETAINED = 3; +export { STALE_GRACE_MS }; + @Injectable() export class StratumV1JobsService { @@ -31,6 +36,7 @@ export class StratumV1JobsService { public blocks: { [id: number]: IJobTemplate } = {}; private lastBlockHeight = 0; + private jobRetentionMs = parseInt(process.env.JOB_RETENTION_MS) || 600000; constructor( private readonly bitcoinRpcService: BitcoinRpcService @@ -109,29 +115,7 @@ export class StratumV1JobsService { } }), tap((data) => { - if (data.blockData.clearJobs) { - this.blocks = {}; - this.jobs = {}; - }else{ - let templatesDeleted = 0; - let jobsDeleted = 0; - const now = new Date().getTime(); - // Delete old templates (5 minutes) - for(const templateId in this.blocks){ - if(now - this.blocks[templateId].blockData.creation > (1000 * 60 * 5)){ - delete this.blocks[templateId]; - templatesDeleted++; - } - } - // Delete old jobs (5 minutes) - for (const jobId in this.jobs) { - if(now - this.jobs[jobId].creation > (1000 * 60 * 5)){ - delete this.jobs[jobId]; - jobsDeleted++; - } - } - //console.log(`Deleted ${templatesDeleted} templates and ${jobsDeleted} jobs.`) - } + this.cleanup(data.blockData.clearJobs); this.blocks[data.blockData.id] = data; }), shareReplay({ refCount: true, bufferSize: 1 }) @@ -162,6 +146,67 @@ export class StratumV1JobsService { return this.blocks[jobTemplateId]; } + public cleanup(clearJobs: boolean, now: number = Date.now()) { + if (clearJobs) { + for (const id in this.blocks) { + const block = this.blocks[id]; + if (block.blockData.retiredAt === undefined) { + block.blockData.retiredAt = now; + } + } + for (const jobId in this.jobs) { + const job = this.jobs[jobId]; + if (job.retiredAt === undefined) { + job.retiredAt = now; + } + } + } + + this.ageEntries( + this.blocks, + now, + entry => entry.blockData.creation, + entry => entry.blockData.retiredAt, + ); + this.ageEntries( + this.jobs, + now, + entry => entry.creation, + entry => entry.retiredAt, + ); + } + + private ageEntries( + map: Record, + now: number, + getCreation: (entry: T) => number, + getRetiredAt: (entry: T) => number | undefined, + ): void { + const ids = Object.keys(map); + if (ids.length <= MIN_RETAINED) { + return; + } + + const candidatesForGC = ids + .slice() + .sort((a, b) => getCreation(map[b]) - getCreation(map[a])) + .slice(MIN_RETAINED); + + for (const id of candidatesForGC) { + const entry = map[id]; + const retiredAt = getRetiredAt(entry); + + if (retiredAt !== undefined && now - retiredAt > this.jobRetentionMs) { + delete map[id]; + continue; + } + + if (retiredAt === undefined && now - getCreation(entry) > this.jobRetentionMs * 2) { + delete map[id]; + } + } + } + public addJob(job: MiningJob) { this.jobs[job.jobId] = job; this.latestJobId++; @@ -171,6 +216,13 @@ export class StratumV1JobsService { return this.jobs[jobId]; } + public classifyJobForShare(job: MiningJob, now: number = Date.now()): 'active' | 'stale-creditable' | 'stale-rejected' { + if (job.retiredAt === undefined) { + return 'active'; + } + return (now - job.retiredAt) <= STALE_GRACE_MS ? 'stale-creditable' : 'stale-rejected'; + } + public getNextTemplateId() { return this.latestJobTemplateId.toString(16); }