mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 17:15:03 -07:00
Compare commits
12
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9878fa42b9 | ||
|
|
8424638081 | ||
|
|
8584ab7efb | ||
|
|
703546719c | ||
|
|
d9fa933825 | ||
|
|
d6d3727119 | ||
|
|
8709de3d83 | ||
|
|
0d15d4fb70 | ||
|
|
5f66d082ea | ||
|
|
93f09ae889 | ||
|
|
216fa74d99 | ||
|
|
355437df3c |
@@ -15,14 +15,17 @@ export class AddressSettingsService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
public async getSettings(address: string, createIfNotFound: boolean) {
|
public async getSettings(address: string, createIfNotFound: boolean) {
|
||||||
const settings = await this.addressSettingsRepository.findOne({ where: { address } });
|
let settings = await this.addressSettingsRepository.findOne({ where: { address } });
|
||||||
if (createIfNotFound == true && settings == null) {
|
if (createIfNotFound === true && settings == null) {
|
||||||
// It's possible to have a race condition here so if we get a PK violation, fetch it
|
await this.addressSettingsRepository
|
||||||
try {
|
.createQueryBuilder()
|
||||||
return await this.createNew(address);
|
.insert()
|
||||||
} catch (e) {
|
.into(AddressSettingsEntity)
|
||||||
return await this.addressSettingsRepository.findOne({ where: { address } });
|
.values({ address })
|
||||||
}
|
.orIgnore()
|
||||||
|
.execute();
|
||||||
|
|
||||||
|
settings = await this.addressSettingsRepository.findOne({ where: { address } });
|
||||||
}
|
}
|
||||||
return settings;
|
return settings;
|
||||||
}
|
}
|
||||||
@@ -81,4 +84,4 @@ export class AddressSettingsService {
|
|||||||
.limit(10)
|
.limit(10)
|
||||||
.execute();
|
.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';
|
import { ValueTransformer } from 'typeorm';
|
||||||
|
|
||||||
export class DateTimeTransformer implements ValueTransformer {
|
const EUROPEAN_LOCALE_REGEX =
|
||||||
to(value: Date): any {
|
/^(\d{1,2})\.(\d{1,2})\.(\d{4}),\s*(\d{1,2}):(\d{2})(?::(\d{2}))?(?:\s*(AM|PM))?$/i;
|
||||||
// Convert the local time to UTC before saving to the database
|
const US_LOCALE_REGEX =
|
||||||
const utcTime = value?.toLocaleString();
|
/^(\d{1,2})\/(\d{1,2})\/(\d{4}),\s*(\d{1,2}):(\d{2})(?::(\d{2}))?\s*(AM|PM)$/i;
|
||||||
return utcTime;
|
|
||||||
|
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 {
|
const matchUs = trimmed.match(US_LOCALE_REGEX);
|
||||||
// Convert the UTC time from the database to the local time zone
|
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;
|
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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -109,4 +109,32 @@ describe('MiningJob', () => {
|
|||||||
|
|
||||||
expect(fastHeader.equals(updatedBlock.toBuffer(true))).toBe(true);
|
expect(fastHeader.equals(updatedBlock.toBuffer(true))).toBe(true);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
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);
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
+41
-22
@@ -3,9 +3,10 @@ import * as bitcoinjs from 'bitcoinjs-lib';
|
|||||||
|
|
||||||
import { IJobTemplate } from '../services/stratum-v1-jobs.service';
|
import { IJobTemplate } from '../services/stratum-v1-jobs.service';
|
||||||
import { eResponseMethod } from './enums/eResponseMethod';
|
import { eResponseMethod } from './enums/eResponseMethod';
|
||||||
import { IMiningNotify } from './stratum-messages/IMiningNotify';
|
|
||||||
import { TOTAL_EXTRANONCE_SIZE_BYTES } from './stratum.constants';
|
import { TOTAL_EXTRANONCE_SIZE_BYTES } from './stratum.constants';
|
||||||
|
|
||||||
|
const MAX_BLOCK_WEIGHT = 4000000;
|
||||||
|
const MAX_SCRIPT_SIZE = 100;
|
||||||
|
|
||||||
interface AddressObject {
|
interface AddressObject {
|
||||||
address: string;
|
address: string;
|
||||||
@@ -19,6 +20,7 @@ export class MiningJob {
|
|||||||
private coinbasePart1Buffer: Buffer;
|
private coinbasePart1Buffer: Buffer;
|
||||||
private coinbasePart2Buffer: Buffer;
|
private coinbasePart2Buffer: Buffer;
|
||||||
private merkleBranchBuffers: Buffer[];
|
private merkleBranchBuffers: Buffer[];
|
||||||
|
private notifyStaticParams: string;
|
||||||
|
|
||||||
public jobTemplateId: string;
|
public jobTemplateId: string;
|
||||||
public networkDifficulty: number;
|
public networkDifficulty: number;
|
||||||
@@ -28,7 +30,8 @@ export class MiningJob {
|
|||||||
private network: bitcoinjs.networks.Network,
|
private network: bitcoinjs.networks.Network,
|
||||||
public jobId: string,
|
public jobId: string,
|
||||||
payoutInformation: AddressObject[],
|
payoutInformation: AddressObject[],
|
||||||
jobTemplate: IJobTemplate
|
jobTemplate: IJobTemplate,
|
||||||
|
poolIdentifier: string = process.env.POOL_IDENTIFIER || 'Public-Pool'
|
||||||
) {
|
) {
|
||||||
|
|
||||||
this.creation = new Date().getTime();
|
this.creation = new Date().getTime();
|
||||||
@@ -45,7 +48,7 @@ export class MiningJob {
|
|||||||
// 32-byte - Commitment hash: Double-SHA256(witness root hash|witness reserved value)
|
// 32-byte - Commitment hash: Double-SHA256(witness root hash|witness reserved value)
|
||||||
|
|
||||||
// 39th byte onwards: Optional data with no consensus meaning
|
// 39th byte onwards: Optional data with no consensus meaning
|
||||||
const extra = Buffer.from('Public-Pool');
|
const extra = Buffer.from(poolIdentifier);
|
||||||
|
|
||||||
// Encode the block height
|
// Encode the block height
|
||||||
// https://github.com/bitcoin/bips/blob/master/bip-0034.mediawiki
|
// https://github.com/bitcoin/bips/blob/master/bip-0034.mediawiki
|
||||||
@@ -58,10 +61,21 @@ export class MiningJob {
|
|||||||
const padding = Buffer.alloc(TOTAL_EXTRANONCE_SIZE_BYTES + (3 - blockHeightEncoded.length), 0)
|
const padding = Buffer.alloc(TOTAL_EXTRANONCE_SIZE_BYTES + (3 - blockHeightEncoded.length), 0)
|
||||||
|
|
||||||
// build the script
|
// 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);
|
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
|
// get the non-witness coinbase tx
|
||||||
//@ts-ignore
|
//@ts-ignore
|
||||||
const serializedCoinbaseTx = this.coinbaseTransaction.__toBuffer().toString('hex');
|
const serializedCoinbaseTx = this.coinbaseTransaction.__toBuffer().toString('hex');
|
||||||
@@ -74,8 +88,30 @@ export class MiningJob {
|
|||||||
this.coinbasePart2 = serializedCoinbaseTx.slice(partOneIndex);
|
this.coinbasePart2 = serializedCoinbaseTx.slice(partOneIndex);
|
||||||
this.coinbasePart1Buffer = Buffer.from(this.coinbasePart1, 'hex');
|
this.coinbasePart1Buffer = Buffer.from(this.coinbasePart1, 'hex');
|
||||||
this.coinbasePart2Buffer = Buffer.from(this.coinbasePart2, 'hex');
|
this.coinbasePart2Buffer = Buffer.from(this.coinbasePart2, 'hex');
|
||||||
|
this.notifyStaticParams = JSON.stringify([
|
||||||
|
this.swapEndianWords(jobTemplate.block.prevHash).toString('hex'),
|
||||||
|
this.coinbasePart1,
|
||||||
|
this.coinbasePart2,
|
||||||
|
jobTemplate.merkle_branch,
|
||||||
|
jobTemplate.block.version.toString(16),
|
||||||
|
jobTemplate.block.bits.toString(16),
|
||||||
|
jobTemplate.block.timestamp.toString(16),
|
||||||
|
jobTemplate.blockData.clearJobs
|
||||||
|
]).slice(1);
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
public getCoinbaseTxHex(): string {
|
||||||
|
//@ts-ignore
|
||||||
|
return this.coinbaseTransaction.__toBuffer().toString('hex');
|
||||||
|
}
|
||||||
|
|
||||||
|
public getCoinbasePrefixBuffer(): Buffer {
|
||||||
|
return Buffer.from(this.coinbasePart1Buffer);
|
||||||
|
}
|
||||||
|
|
||||||
|
public getCoinbaseSuffixBuffer(): Buffer {
|
||||||
|
return Buffer.from(this.coinbasePart2Buffer);
|
||||||
}
|
}
|
||||||
|
|
||||||
public cloneCoinbaseTransaction(): bitcoinjs.Transaction {
|
public cloneCoinbaseTransaction(): bitcoinjs.Transaction {
|
||||||
@@ -209,24 +245,7 @@ export class MiningJob {
|
|||||||
}
|
}
|
||||||
|
|
||||||
public response(jobTemplate: IJobTemplate): string {
|
public response(jobTemplate: IJobTemplate): string {
|
||||||
|
return `{"id":null,"method":"${eResponseMethod.MINING_NOTIFY}","params":[${JSON.stringify(this.jobId)},${this.notifyStaticParams}}\n`;
|
||||||
const job: IMiningNotify = {
|
|
||||||
id: null,
|
|
||||||
method: eResponseMethod.MINING_NOTIFY,
|
|
||||||
params: [
|
|
||||||
this.jobId,
|
|
||||||
this.swapEndianWords(jobTemplate.block.prevHash).toString('hex'),
|
|
||||||
this.coinbasePart1,
|
|
||||||
this.coinbasePart2,
|
|
||||||
jobTemplate.merkle_branch,
|
|
||||||
jobTemplate.block.version.toString(16),
|
|
||||||
jobTemplate.block.bits.toString(16),
|
|
||||||
jobTemplate.block.timestamp.toString(16),
|
|
||||||
jobTemplate.blockData.clearJobs
|
|
||||||
]
|
|
||||||
};
|
|
||||||
|
|
||||||
return JSON.stringify(job) + '\n';
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ import { NotificationService } from '../services/notification.service';
|
|||||||
import { StratumV1JobsService } from '../services/stratum-v1-jobs.service';
|
import { StratumV1JobsService } from '../services/stratum-v1-jobs.service';
|
||||||
import { IBlockTemplate } from './bitcoin-rpc/IBlockTemplate';
|
import { IBlockTemplate } from './bitcoin-rpc/IBlockTemplate';
|
||||||
import { MiningJob } from './MiningJob';
|
import { MiningJob } from './MiningJob';
|
||||||
import { StratumV1Client } from './StratumV1Client';
|
import { effectiveJobDifficulty, StratumV1Client } from './StratumV1Client';
|
||||||
import { MiningSubmitMessage } from './stratum-messages/MiningSubmitMessage';
|
import { MiningSubmitMessage } from './stratum-messages/MiningSubmitMessage';
|
||||||
|
|
||||||
|
|
||||||
@@ -199,6 +199,27 @@ describe('StratumV1Client', () => {
|
|||||||
expect(socket.on).toHaveBeenCalled();
|
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 () => {
|
it('should respond to mining.subscribe', async () => {
|
||||||
jest.spyOn(socket, 'write').mockImplementation((data) => true);
|
jest.spyOn(socket, 'write').mockImplementation((data) => true);
|
||||||
|
|
||||||
@@ -576,12 +597,14 @@ describe('StratumV1Client', () => {
|
|||||||
|
|
||||||
it('should submit and persist found blocks', async () => {
|
it('should submit and persist found blocks', async () => {
|
||||||
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
|
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({
|
jest.spyOn(client as any, 'calculateDifficulty').mockReturnValue({
|
||||||
submissionDifficulty: Number.MAX_SAFE_INTEGER,
|
submissionDifficulty: Number.MAX_SAFE_INTEGER,
|
||||||
submissionHash: 'block-share'
|
submissionHash: 'block-share'
|
||||||
});
|
});
|
||||||
const addressSettings = moduleRef.get<AddressSettingsService>(AddressSettingsService);
|
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(MockRecording1.MINING_SUBSCRIBE);
|
||||||
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
|
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
|
||||||
@@ -605,6 +628,43 @@ describe('StratumV1Client', () => {
|
|||||||
expect((client as any).write).lastCalledWith(`{"id":5,"error":null,"result":true}\n`);
|
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);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
@@ -32,6 +32,22 @@ const TRUE_DIFF_ONE = 2.695953529101131e67;
|
|||||||
const BLOCKED_USER_AGENT_LOG_INTERVAL_MS = 60 * 1000;
|
const BLOCKED_USER_AGENT_LOG_INTERVAL_MS = 60 * 1000;
|
||||||
const VALIDATION_ERROR_LOG_INTERVAL_MS = 60 * 1000;
|
const VALIDATION_ERROR_LOG_INTERVAL_MS = 60 * 1000;
|
||||||
|
|
||||||
|
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 {
|
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 }>();
|
||||||
@@ -47,6 +63,8 @@ export class StratumV1Client {
|
|||||||
private stratumInitialized = false;
|
private stratumInitialized = false;
|
||||||
private usedSuggestedDifficulty = false;
|
private usedSuggestedDifficulty = false;
|
||||||
private sessionDifficulty: number = 100000;
|
private sessionDifficulty: number = 100000;
|
||||||
|
private oldSessionDifficulty: number = 100000;
|
||||||
|
private diffChangeJobId: number | null = null;
|
||||||
|
|
||||||
private clientEntity: ClientEntity;
|
private clientEntity: ClientEntity;
|
||||||
private creatingEntity: Promise<void>;
|
private creatingEntity: Promise<void>;
|
||||||
@@ -237,6 +255,7 @@ export class StratumV1Client {
|
|||||||
this.clientAuthorization = authorizationMessage;
|
this.clientAuthorization = authorizationMessage;
|
||||||
if (this.clientSuggestedDifficulty == null && this.clientAuthorization.startingDiff != null && this.clientAuthorization.startingDiff > this.sessionDifficulty) {
|
if (this.clientSuggestedDifficulty == null && this.clientAuthorization.startingDiff != null && this.clientAuthorization.startingDiff > this.sessionDifficulty) {
|
||||||
this.sessionDifficulty = this.clientAuthorization.startingDiff;
|
this.sessionDifficulty = this.clientAuthorization.startingDiff;
|
||||||
|
this.oldSessionDifficulty = this.sessionDifficulty;
|
||||||
}
|
}
|
||||||
const success = await this.write(JSON.stringify(this.clientAuthorization.response()) + '\n');
|
const success = await this.write(JSON.stringify(this.clientAuthorization.response()) + '\n');
|
||||||
if (!success) {
|
if (!success) {
|
||||||
@@ -278,6 +297,7 @@ export class StratumV1Client {
|
|||||||
|
|
||||||
this.clientSuggestedDifficulty = suggestDifficultyMessage;
|
this.clientSuggestedDifficulty = suggestDifficultyMessage;
|
||||||
this.sessionDifficulty = suggestDifficultyMessage.suggestedDifficulty;
|
this.sessionDifficulty = suggestDifficultyMessage.suggestedDifficulty;
|
||||||
|
this.oldSessionDifficulty = this.sessionDifficulty;
|
||||||
const success = await this.write(JSON.stringify(this.clientSuggestedDifficulty.response(this.sessionDifficulty)) + '\n');
|
const success = await this.write(JSON.stringify(this.clientSuggestedDifficulty.response(this.sessionDifficulty)) + '\n');
|
||||||
if (!success) {
|
if (!success) {
|
||||||
return;
|
return;
|
||||||
@@ -375,6 +395,7 @@ export class StratumV1Client {
|
|||||||
switch (this.clientSubscription.userAgent) {
|
switch (this.clientSubscription.userAgent) {
|
||||||
case 'cpuminer': {
|
case 'cpuminer': {
|
||||||
this.sessionDifficulty = 0.1;
|
this.sessionDifficulty = 0.1;
|
||||||
|
this.oldSessionDifficulty = this.sessionDifficulty;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -457,7 +478,8 @@ export class StratumV1Client {
|
|||||||
network,
|
network,
|
||||||
this.stratumV1JobsService.getNextId(),
|
this.stratumV1JobsService.getNextId(),
|
||||||
payoutInformation,
|
payoutInformation,
|
||||||
jobTemplate
|
jobTemplate,
|
||||||
|
this.configService.get('POOL_IDENTIFIER') || 'Public-Pool'
|
||||||
);
|
);
|
||||||
|
|
||||||
this.stratumV1JobsService.addJob(job);
|
this.stratumV1JobsService.addJob(job);
|
||||||
@@ -564,9 +586,15 @@ export class StratumV1Client {
|
|||||||
const { submissionDifficulty } = this.calculateDifficulty(header);
|
const { submissionDifficulty } = this.calculateDifficulty(header);
|
||||||
|
|
||||||
//console.log(`DIFF: ${submissionDifficulty} of ${this.sessionDifficulty} from ${this.clientAuthorization.worker + '.' + this.extraNonceAndSessionId}`);
|
//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) {
|
||||||
const success = await this.write(JSON.stringify(submission.response()) + '\n');
|
const success = await this.write(JSON.stringify(submission.response()) + '\n');
|
||||||
if (!success) {
|
if (!success) {
|
||||||
return false;
|
return false;
|
||||||
@@ -584,23 +612,24 @@ export class StratumV1Client {
|
|||||||
);
|
);
|
||||||
const blockHex = updatedJobBlock.toHex(false);
|
const blockHex = updatedJobBlock.toHex(false);
|
||||||
const result = await this.bitcoinRpcService.SUBMIT_BLOCK(blockHex);
|
const result = await this.bitcoinRpcService.SUBMIT_BLOCK(blockHex);
|
||||||
await this.blocksService.save({
|
if (result === 'SUCCESS!') {
|
||||||
height: jobTemplate.blockData.height,
|
await this.blocksService.save({
|
||||||
minerAddress: this.clientAuthorization.address,
|
height: jobTemplate.blockData.height,
|
||||||
worker: this.clientAuthorization.worker,
|
minerAddress: this.clientAuthorization.address,
|
||||||
sessionId: this.extraNonceAndSessionId,
|
worker: this.clientAuthorization.worker,
|
||||||
blockData: blockHex
|
sessionId: this.extraNonceAndSessionId,
|
||||||
});
|
blockData: blockHex
|
||||||
|
});
|
||||||
|
|
||||||
await this.notificationService.notifySubscribersBlockFound(this.clientAuthorization.address, jobTemplate.blockData.height, updatedJobBlock, result);
|
await this.notificationService.notifySubscribersBlockFound(this.clientAuthorization.address, jobTemplate.blockData.height, updatedJobBlock, result);
|
||||||
//success
|
|
||||||
if (result == null) {
|
|
||||||
await this.addressSettingsService.resetBestDifficultyAndShares();
|
await this.addressSettingsService.resetBestDifficultyAndShares();
|
||||||
|
} else {
|
||||||
|
console.warn(`[Block submit rejected at height ${jobTemplate.blockData.height}]: ${result}`);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
await this.ensureClientEntity();
|
await this.ensureClientEntity();
|
||||||
try {
|
try {
|
||||||
await this.statistics.addShares(this.clientEntity, this.sessionDifficulty);
|
await this.statistics.addShares(this.clientEntity, effectiveDiff);
|
||||||
const now = new Date();
|
const now = new Date();
|
||||||
// only update every minute
|
// only update every minute
|
||||||
//if (this.clientEntity.updatedAt == null || now.getTime() - this.clientEntity.updatedAt.getTime() > 1000 * 60) {
|
//if (this.clientEntity.updatedAt == null || now.getTime() - this.clientEntity.updatedAt.getTime() > 1000 * 60) {
|
||||||
@@ -647,6 +676,11 @@ export class StratumV1Client {
|
|||||||
|
|
||||||
if (targetDiff != this.sessionDifficulty) {
|
if (targetDiff != this.sessionDifficulty) {
|
||||||
//console.log(`Adjusting ${this.extraNonceAndSessionId} difficulty from ${this.sessionDifficulty} to ${targetDiff}`);
|
//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;
|
this.sessionDifficulty = targetDiff;
|
||||||
|
|
||||||
const data = JSON.stringify({
|
const data = JSON.stringify({
|
||||||
@@ -656,7 +690,10 @@ export class StratumV1Client {
|
|||||||
}) + '\n';
|
}) + '\n';
|
||||||
|
|
||||||
|
|
||||||
await this.socket.write(data);
|
const success = await this.write(data);
|
||||||
|
if (!success) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
const jobTemplate = await firstValueFrom(this.stratumV1JobsService.newMiningJob$);
|
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
|
// 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();
|
expect(statistics.getSuggestedDifficulty(64)).toBeNull();
|
||||||
});
|
});
|
||||||
|
|
||||||
it('should lower difficulty when a miner has not submitted shares for several minutes', () => {
|
it('should lower difficulty when a miner has not submitted shares for more than a minute', () => {
|
||||||
jest.setSystemTime(new Date('2026-05-06T12:06:00Z'));
|
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 () => {
|
it('should increase difficulty for rapid submissions', async () => {
|
||||||
|
|||||||
@@ -104,8 +104,8 @@ export class StratumV1ClientStatistics {
|
|||||||
|
|
||||||
// miner hasn't submitted shares in one minute
|
// miner hasn't submitted shares in one minute
|
||||||
if (this.submissionCache.length < 5) {
|
if (this.submissionCache.length < 5) {
|
||||||
if ((new Date().getTime() - this.submissionCacheStart.getTime()) / 5000 > 60) {
|
if ((new Date().getTime() - this.submissionCacheStart.getTime()) / 1000 > 60) {
|
||||||
return this.nearestPowerOfTwo(clientDifficulty / 6);
|
return this.nearestDifficultyStep(clientDifficulty / 6);
|
||||||
} else {
|
} else {
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
@@ -116,39 +116,45 @@ 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 (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(difficultyPerSecond) || !Number.isFinite(targetDifficulty)) {
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
if ((clientDifficulty * 2) < targetDifficulty || (clientDifficulty / 2) > targetDifficulty) {
|
if ((clientDifficulty * 2) < targetDifficulty || (clientDifficulty / 2) > targetDifficulty) {
|
||||||
return this.nearestPowerOfTwo(targetDifficulty)
|
return this.nearestDifficultyStep(targetDifficulty)
|
||||||
}
|
}
|
||||||
|
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
private nearestPowerOfTwo(val): number {
|
private nearestDifficultyStep(val: number): number {
|
||||||
if (val === 0) {
|
if (val === 0) {
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
if (val < MIN_DIFF) {
|
if (val < MIN_DIFF) {
|
||||||
return MIN_DIFF;
|
return MIN_DIFF;
|
||||||
}
|
}
|
||||||
let x = val | (val >> 1);
|
|
||||||
x = x | (x >> 2);
|
const exponent = Math.floor(Math.log2(val));
|
||||||
x = x | (x >> 4);
|
const lower = 2 ** exponent;
|
||||||
x = x | (x >> 8);
|
const middle = lower + lower / 2;
|
||||||
x = x | (x >> 16);
|
const upper = lower * 2;
|
||||||
x = x | (x >> 32);
|
|
||||||
const res = x - (x >> 1);
|
const distances = [
|
||||||
if (res == 0 && val * 100 < MIN_DIFF) {
|
{ value: lower, diff: Math.abs(val - lower) },
|
||||||
return MIN_DIFF;
|
{ value: middle, diff: Math.abs(val - middle) },
|
||||||
}
|
{ value: upper, diff: Math.abs(val - upper) },
|
||||||
if (res == 0) {
|
];
|
||||||
return this.nearestPowerOfTwo(val * 100) / 100;
|
|
||||||
}
|
distances.sort((a, b) => a.diff - b.diff);
|
||||||
return res;
|
return distances[0].value;
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -58,6 +58,29 @@ describe('MiningSubmitMessage', () => {
|
|||||||
|
|
||||||
expect(errors.some(error => error.property === 'extraNonce2')).toBe(true);
|
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{
|
public hash(): string{
|
||||||
const buffer = Buffer.from(this.versionMask + this.nonce + this.extraNonce2 + this.ntime + this.jobId);
|
const canonical = JSON.stringify({
|
||||||
return bitcoinjs.crypto.hash256(buffer).toString('base64');
|
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 { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
|
||||||
import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo';
|
import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo';
|
||||||
import * as PGPubsub from 'pg-pubsub';
|
import * as PGPubsub from 'pg-pubsub';
|
||||||
|
import * as fs from 'node:fs';
|
||||||
|
|
||||||
@Injectable()
|
@Injectable()
|
||||||
export class BitcoinRpcService implements OnModuleInit {
|
export class BitcoinRpcService implements OnModuleInit {
|
||||||
@@ -35,10 +36,17 @@ export class BitcoinRpcService implements OnModuleInit {
|
|||||||
|
|
||||||
|
|
||||||
const url = this.configService.get('BITCOIN_RPC_URL');
|
const url = this.configService.get('BITCOIN_RPC_URL');
|
||||||
const user = this.configService.get('BITCOIN_RPC_USER');
|
let user = this.configService.get('BITCOIN_RPC_USER');
|
||||||
const pass = this.configService.get('BITCOIN_RPC_PASSWORD');
|
let pass = this.configService.get('BITCOIN_RPC_PASSWORD');
|
||||||
const port = parseInt(this.configService.get('BITCOIN_RPC_PORT'));
|
const port = parseInt(this.configService.get('BITCOIN_RPC_PORT'));
|
||||||
const timeout = parseInt(this.configService.get('BITCOIN_RPC_TIMEOUT'));
|
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;
|
||||||
|
}
|
||||||
|
|
||||||
const baseURL = this.buildRpcUrl(url, port);
|
const baseURL = this.buildRpcUrl(url, port);
|
||||||
this.client = axios.create({
|
this.client = axios.create({
|
||||||
|
|||||||
@@ -109,29 +109,7 @@ export class StratumV1JobsService {
|
|||||||
}
|
}
|
||||||
}),
|
}),
|
||||||
tap((data) => {
|
tap((data) => {
|
||||||
if (data.blockData.clearJobs) {
|
this.cleanup(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.blocks[data.blockData.id] = data;
|
this.blocks[data.blockData.id] = data;
|
||||||
}),
|
}),
|
||||||
shareReplay({ refCount: true, bufferSize: 1 })
|
shareReplay({ refCount: true, bufferSize: 1 })
|
||||||
@@ -162,6 +140,32 @@ export class StratumV1JobsService {
|
|||||||
return this.blocks[jobTemplateId];
|
return this.blocks[jobTemplateId];
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public cleanup(clearJobs: boolean, now: number = Date.now()) {
|
||||||
|
if (clearJobs) {
|
||||||
|
this.blocks = {};
|
||||||
|
this.jobs = {};
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
let templatesDeleted = 0;
|
||||||
|
let jobsDeleted = 0;
|
||||||
|
|
||||||
|
for (const templateId in this.blocks) {
|
||||||
|
if (now - this.blocks[templateId].blockData.creation > (1000 * 60 * 5)) {
|
||||||
|
delete this.blocks[templateId];
|
||||||
|
templatesDeleted++;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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.`)
|
||||||
|
}
|
||||||
|
|
||||||
public addJob(job: MiningJob) {
|
public addJob(job: MiningJob) {
|
||||||
this.jobs[job.jobId] = job;
|
this.jobs[job.jobId] = job;
|
||||||
this.latestJobId++;
|
this.latestJobId++;
|
||||||
|
|||||||
@@ -5,12 +5,10 @@ describe('StratumV1Service', () => {
|
|||||||
const originalStratumPorts = process.env.STRATUM_PORTS;
|
const originalStratumPorts = process.env.STRATUM_PORTS;
|
||||||
const originalStratumSecure = process.env.STRATUM_SECURE;
|
const originalStratumSecure = process.env.STRATUM_SECURE;
|
||||||
const originalSecureStratumPorts = process.env.SECURE_STRATUM_PORTS;
|
const originalSecureStratumPorts = process.env.SECURE_STRATUM_PORTS;
|
||||||
const originalBackpressureEnabled = process.env.STRATUM_BACKPRESSURE_ENABLED;
|
|
||||||
|
|
||||||
let service: StratumV1Service;
|
let service: StratumV1Service;
|
||||||
let clientService;
|
let clientService;
|
||||||
let consoleLogSpy: jest.SpyInstance;
|
let consoleLogSpy: jest.SpyInstance;
|
||||||
let consoleWarnSpy: jest.SpyInstance;
|
|
||||||
|
|
||||||
beforeEach(() => {
|
beforeEach(() => {
|
||||||
jest.useFakeTimers();
|
jest.useFakeTimers();
|
||||||
@@ -28,7 +26,6 @@ describe('StratumV1Service', () => {
|
|||||||
{} as any
|
{} as any
|
||||||
);
|
);
|
||||||
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
|
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
|
||||||
consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined);
|
|
||||||
});
|
});
|
||||||
|
|
||||||
afterEach(() => {
|
afterEach(() => {
|
||||||
@@ -36,9 +33,7 @@ describe('StratumV1Service', () => {
|
|||||||
restoreEnv('STRATUM_PORTS', originalStratumPorts);
|
restoreEnv('STRATUM_PORTS', originalStratumPorts);
|
||||||
restoreEnv('STRATUM_SECURE', originalStratumSecure);
|
restoreEnv('STRATUM_SECURE', originalStratumSecure);
|
||||||
restoreEnv('SECURE_STRATUM_PORTS', originalSecureStratumPorts);
|
restoreEnv('SECURE_STRATUM_PORTS', originalSecureStratumPorts);
|
||||||
restoreEnv('STRATUM_BACKPRESSURE_ENABLED', originalBackpressureEnabled);
|
|
||||||
consoleLogSpy.mockRestore();
|
consoleLogSpy.mockRestore();
|
||||||
consoleWarnSpy.mockRestore();
|
|
||||||
jest.useRealTimers();
|
jest.useRealTimers();
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -73,52 +68,6 @@ describe('StratumV1Service', () => {
|
|||||||
expect(startSecureSocketServerSpy).toHaveBeenCalledWith(4333);
|
expect(startSecureSocketServerSpy).toHaveBeenCalledWith(4333);
|
||||||
});
|
});
|
||||||
|
|
||||||
it('should pause listeners when worker backpressure is high', () => {
|
|
||||||
const close = jest.fn((callback?: (error?: Error) => void) => callback?.());
|
|
||||||
(service as any).listeners.push({
|
|
||||||
port: 3333,
|
|
||||||
secure: false,
|
|
||||||
server: { close },
|
|
||||||
paused: false
|
|
||||||
});
|
|
||||||
jest.spyOn(service as any, 'getEventLoopP95Ms').mockReturnValue(5000);
|
|
||||||
jest.spyOn(service as any, 'getBackpressureEventLoopP95Ms').mockReturnValue(2000);
|
|
||||||
|
|
||||||
(service as any).checkBackpressure();
|
|
||||||
|
|
||||||
expect(close).toHaveBeenCalled();
|
|
||||||
expect((service as any).listeners[0].paused).toBe(true);
|
|
||||||
expect((service as any).listeners[0].server).toBeNull();
|
|
||||||
expect(consoleWarnSpy).toHaveBeenCalledWith(expect.stringContaining('Pausing Stratum accepts'));
|
|
||||||
});
|
|
||||||
|
|
||||||
it('should resume listeners after consecutive healthy backpressure checks', () => {
|
|
||||||
(service as any).listeners.push({
|
|
||||||
port: 3333,
|
|
||||||
secure: false,
|
|
||||||
server: null,
|
|
||||||
paused: true
|
|
||||||
});
|
|
||||||
jest.spyOn(service as any, 'getEventLoopP95Ms').mockReturnValue(50);
|
|
||||||
jest.spyOn(service as any, 'getBackpressureEventLoopP95Ms').mockReturnValue(2000);
|
|
||||||
jest.spyOn(service as any, 'getBackpressureResumeEventLoopP95Ms').mockReturnValue(250);
|
|
||||||
jest.spyOn(service as any, 'getBackpressureResumeRssMb').mockReturnValue(Number.MAX_SAFE_INTEGER);
|
|
||||||
jest.spyOn(service as any, 'getBackpressureHealthyChecks').mockReturnValue(2);
|
|
||||||
const listenSpy = jest.spyOn(service as any, 'listen').mockImplementation((listener: any) => {
|
|
||||||
listener.server = {};
|
|
||||||
listener.paused = false;
|
|
||||||
});
|
|
||||||
|
|
||||||
(service as any).checkBackpressure();
|
|
||||||
expect(listenSpy).not.toHaveBeenCalled();
|
|
||||||
|
|
||||||
(service as any).checkBackpressure();
|
|
||||||
|
|
||||||
expect(listenSpy).toHaveBeenCalledWith((service as any).listeners[0]);
|
|
||||||
expect((service as any).listeners[0].paused).toBe(false);
|
|
||||||
expect(consoleWarnSpy).toHaveBeenCalledWith(expect.stringContaining('Resuming Stratum accepts'));
|
|
||||||
});
|
|
||||||
|
|
||||||
function restoreEnv(key: string, value: string | undefined) {
|
function restoreEnv(key: string, value: string | undefined) {
|
||||||
if (value == null) {
|
if (value == null) {
|
||||||
delete process.env[key];
|
delete process.env[key];
|
||||||
|
|||||||
@@ -1,7 +1,6 @@
|
|||||||
import { Injectable, OnModuleInit } from '@nestjs/common';
|
import { Injectable, OnModuleInit } from '@nestjs/common';
|
||||||
import { ConfigService } from '@nestjs/config';
|
import { ConfigService } from '@nestjs/config';
|
||||||
import { Server, Socket } from 'net';
|
import { Server, Socket } from 'net';
|
||||||
import { monitorEventLoopDelay } from 'perf_hooks';
|
|
||||||
|
|
||||||
import { StratumV1Client } from '../models/StratumV1Client';
|
import { StratumV1Client } from '../models/StratumV1Client';
|
||||||
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
|
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
|
||||||
@@ -16,20 +15,6 @@ import { readFileSync } from 'fs';
|
|||||||
import { TlsOptions, TLSSocket, createServer } from 'tls';
|
import { TlsOptions, TLSSocket, createServer } from 'tls';
|
||||||
import * as path from 'path';
|
import * as path from 'path';
|
||||||
|
|
||||||
interface StratumListenerState {
|
|
||||||
port: number;
|
|
||||||
secure: boolean;
|
|
||||||
server: Server | null;
|
|
||||||
paused: boolean;
|
|
||||||
}
|
|
||||||
|
|
||||||
const DEFAULT_BACKPRESSURE_CHECK_INTERVAL_MS = 5000;
|
|
||||||
const DEFAULT_BACKPRESSURE_EVENT_LOOP_P95_MS = 2000;
|
|
||||||
const DEFAULT_BACKPRESSURE_EVENT_LOOP_RESUME_P95_MS = 250;
|
|
||||||
const DEFAULT_BACKPRESSURE_RSS_MB = 2500;
|
|
||||||
const DEFAULT_BACKPRESSURE_RESUME_RSS_MB = 2000;
|
|
||||||
const DEFAULT_BACKPRESSURE_HEALTHY_CHECKS = 3;
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
@Injectable()
|
@Injectable()
|
||||||
@@ -39,10 +24,6 @@ export class StratumV1Service implements OnModuleInit {
|
|||||||
private emptySocket = 0;
|
private emptySocket = 0;
|
||||||
private normalClosure = 0;
|
private normalClosure = 0;
|
||||||
private errorClosure = 0;
|
private errorClosure = 0;
|
||||||
private readonly listeners: StratumListenerState[] = [];
|
|
||||||
private readonly eventLoopDelay = monitorEventLoopDelay({ resolution: 20 });
|
|
||||||
private backpressureMonitor: NodeJS.Timeout | null = null;
|
|
||||||
private healthyBackpressureChecks = 0;
|
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
private readonly bitcoinRpcService: BitcoinRpcService,
|
private readonly bitcoinRpcService: BitcoinRpcService,
|
||||||
@@ -85,22 +66,9 @@ export class StratumV1Service implements OnModuleInit {
|
|||||||
this.errorClosure = 0;
|
this.errorClosure = 0;
|
||||||
}, 1000 * 60);
|
}, 1000 * 60);
|
||||||
|
|
||||||
this.startBackpressureMonitor();
|
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private startSocketServer(port: number) {
|
private startSocketServer(port: number) {
|
||||||
const listener: StratumListenerState = {
|
|
||||||
port,
|
|
||||||
secure: false,
|
|
||||||
server: null,
|
|
||||||
paused: false
|
|
||||||
};
|
|
||||||
this.listeners.push(listener);
|
|
||||||
this.listen(listener);
|
|
||||||
}
|
|
||||||
|
|
||||||
private createSocketServer(): Server {
|
|
||||||
const server = new Server(async (socket: Socket) => {
|
const server = new Server(async (socket: Socket) => {
|
||||||
// Set 15-minute timeout
|
// Set 15-minute timeout
|
||||||
socket.setTimeout(1000 * 60 * 15);
|
socket.setTimeout(1000 * 60 * 15);
|
||||||
@@ -163,21 +131,13 @@ export class StratumV1Service implements OnModuleInit {
|
|||||||
console.error(`Server error: ${err.message}`);
|
console.error(`Server error: ${err.message}`);
|
||||||
});
|
});
|
||||||
|
|
||||||
return server;
|
server.listen(port, () => {
|
||||||
|
console.log(`Stratum server is listening on port ${port}`);
|
||||||
|
});
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private startSecureSocketServer(port: number) {
|
private startSecureSocketServer(port: number) {
|
||||||
const listener: StratumListenerState = {
|
|
||||||
port,
|
|
||||||
secure: true,
|
|
||||||
server: null,
|
|
||||||
paused: false
|
|
||||||
};
|
|
||||||
this.listeners.push(listener);
|
|
||||||
this.listen(listener);
|
|
||||||
}
|
|
||||||
|
|
||||||
private createSecureSocketServer(): Server {
|
|
||||||
|
|
||||||
const currentDirectory = process.cwd();
|
const currentDirectory = process.cwd();
|
||||||
const keyPath = path.join(currentDirectory, 'secrets', 'key.pem');
|
const keyPath = path.join(currentDirectory, 'secrets', 'key.pem');
|
||||||
@@ -243,140 +203,9 @@ export class StratumV1Service implements OnModuleInit {
|
|||||||
console.error(`Server error: ${err.message}`);
|
console.error(`Server error: ${err.message}`);
|
||||||
});
|
});
|
||||||
|
|
||||||
return server;
|
server.listen(port, () => {
|
||||||
|
console.log(`Stratum TLS server is listening on port ${port}`);
|
||||||
}
|
|
||||||
|
|
||||||
private listen(listener: StratumListenerState) {
|
|
||||||
if (listener.server != null) {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
const server = listener.secure ? this.createSecureSocketServer() : this.createSocketServer();
|
|
||||||
listener.server = server;
|
|
||||||
listener.paused = false;
|
|
||||||
|
|
||||||
server.listen(listener.port, () => {
|
|
||||||
console.log(`${listener.secure ? 'Stratum TLS' : 'Stratum'} server is listening on port ${listener.port}`);
|
|
||||||
});
|
});
|
||||||
}
|
|
||||||
|
|
||||||
private startBackpressureMonitor() {
|
|
||||||
if (this.isBackpressureDisabled() || this.backpressureMonitor != null) {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
this.eventLoopDelay.enable();
|
|
||||||
this.backpressureMonitor = setInterval(() => {
|
|
||||||
this.checkBackpressure();
|
|
||||||
}, this.getBackpressureCheckIntervalMs());
|
|
||||||
}
|
|
||||||
|
|
||||||
private checkBackpressure() {
|
|
||||||
const eventLoopP95Ms = this.getEventLoopP95Ms();
|
|
||||||
const rssMb = Math.round(process.memoryUsage().rss / 1024 / 1024);
|
|
||||||
const overloaded = eventLoopP95Ms >= this.getBackpressureEventLoopP95Ms()
|
|
||||||
|| rssMb >= this.getBackpressureRssMb();
|
|
||||||
const paused = this.listeners.some(listener => listener.paused);
|
|
||||||
|
|
||||||
if (overloaded) {
|
|
||||||
this.healthyBackpressureChecks = 0;
|
|
||||||
if (!paused) {
|
|
||||||
this.pauseAccepting(eventLoopP95Ms, rssMb);
|
|
||||||
}
|
|
||||||
this.eventLoopDelay.reset();
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (!paused) {
|
|
||||||
this.eventLoopDelay.reset();
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
const healthy = eventLoopP95Ms <= this.getBackpressureResumeEventLoopP95Ms()
|
|
||||||
&& rssMb <= this.getBackpressureResumeRssMb();
|
|
||||||
if (!healthy) {
|
|
||||||
this.healthyBackpressureChecks = 0;
|
|
||||||
this.eventLoopDelay.reset();
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
this.healthyBackpressureChecks++;
|
|
||||||
if (this.healthyBackpressureChecks >= this.getBackpressureHealthyChecks()) {
|
|
||||||
this.resumeAccepting(eventLoopP95Ms, rssMb);
|
|
||||||
this.healthyBackpressureChecks = 0;
|
|
||||||
}
|
|
||||||
|
|
||||||
this.eventLoopDelay.reset();
|
|
||||||
}
|
|
||||||
|
|
||||||
private pauseAccepting(eventLoopP95Ms: number, rssMb: number) {
|
|
||||||
console.warn(`Pausing Stratum accepts: eventLoopP95Ms=${eventLoopP95Ms}, rssMb=${rssMb}`);
|
|
||||||
for (const listener of this.listeners) {
|
|
||||||
if (listener.paused || listener.server == null) {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
const server = listener.server;
|
|
||||||
listener.server = null;
|
|
||||||
listener.paused = true;
|
|
||||||
server.close((error) => {
|
|
||||||
if (error != null) {
|
|
||||||
console.error(`Error while pausing Stratum listener on port ${listener.port}: ${error.message}`);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private resumeAccepting(eventLoopP95Ms: number, rssMb: number) {
|
|
||||||
console.warn(`Resuming Stratum accepts: eventLoopP95Ms=${eventLoopP95Ms}, rssMb=${rssMb}`);
|
|
||||||
for (const listener of this.listeners) {
|
|
||||||
if (!listener.paused || listener.server != null) {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
this.listen(listener);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private getEventLoopP95Ms() {
|
|
||||||
return Math.round(this.eventLoopDelay.percentile(95) / 1e6);
|
|
||||||
}
|
|
||||||
|
|
||||||
private isBackpressureDisabled() {
|
|
||||||
return process.env.STRATUM_BACKPRESSURE_ENABLED?.toLowerCase() === 'false';
|
|
||||||
}
|
|
||||||
|
|
||||||
private getBackpressureCheckIntervalMs() {
|
|
||||||
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_CHECK_INTERVAL_MS', DEFAULT_BACKPRESSURE_CHECK_INTERVAL_MS);
|
|
||||||
}
|
|
||||||
|
|
||||||
private getBackpressureEventLoopP95Ms() {
|
|
||||||
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_EVENT_LOOP_P95_MS', DEFAULT_BACKPRESSURE_EVENT_LOOP_P95_MS);
|
|
||||||
}
|
|
||||||
|
|
||||||
private getBackpressureResumeEventLoopP95Ms() {
|
|
||||||
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_EVENT_LOOP_RESUME_P95_MS', DEFAULT_BACKPRESSURE_EVENT_LOOP_RESUME_P95_MS);
|
|
||||||
}
|
|
||||||
|
|
||||||
private getBackpressureRssMb() {
|
|
||||||
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_RSS_MB', DEFAULT_BACKPRESSURE_RSS_MB);
|
|
||||||
}
|
|
||||||
|
|
||||||
private getBackpressureResumeRssMb() {
|
|
||||||
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_RESUME_RSS_MB', DEFAULT_BACKPRESSURE_RESUME_RSS_MB);
|
|
||||||
}
|
|
||||||
|
|
||||||
private getBackpressureHealthyChecks() {
|
|
||||||
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_HEALTHY_CHECKS', DEFAULT_BACKPRESSURE_HEALTHY_CHECKS);
|
|
||||||
}
|
|
||||||
|
|
||||||
private getPositiveIntegerEnv(key: string, fallback: number) {
|
|
||||||
const configured = parseInt(process.env[key], 10);
|
|
||||||
if (Number.isFinite(configured) && configured > 0) {
|
|
||||||
return configured;
|
|
||||||
}
|
|
||||||
return fallback;
|
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user