mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 09:05:06 -07:00
misc improvements from blitzpool
This commit is contained in:
@@ -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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
});
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
|
||||
+35
-3
@@ -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);
|
||||
|
||||
@@ -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<IBlockTemplate> = 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>(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>(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>(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);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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<void>;
|
||||
@@ -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
|
||||
|
||||
@@ -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 () => {
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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());
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
|
||||
@@ -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');
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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 });
|
||||
|
||||
|
||||
@@ -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');
|
||||
|
||||
|
||||
@@ -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<T>(
|
||||
map: Record<string, T>,
|
||||
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);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user