11 Commits
Author SHA1 Message Date
Ben 8424638081 Revert "improve stratum connection and broadcast performance"
This reverts commit d6d3727119.
2026-05-06 21:37:57 -04:00
Ben 8584ab7efb Revert "optimize stratum broadcast fanout"
This reverts commit d9fa933825.
2026-05-06 21:37:43 -04:00
Ben 703546719c optimizations 2026-05-06 21:24:04 -04:00
Ben d9fa933825 optimize stratum broadcast fanout 2026-05-06 20:58:08 -04:00
Ben d6d3727119 improve stratum connection and broadcast performance 2026-05-06 20:38:42 -04:00
Ben 8709de3d83 reduce stratum rejection overhead 2026-05-06 20:26:00 -04:00
Ben 0d15d4fb70 reduce overhead 2026-05-06 20:15:27 -04:00
Ben 5f66d082ea dependency updates 2026-05-06 19:59:58 -04:00
Ben 93f09ae889 optimize share submission path 2026-05-06 19:45:32 -04:00
Ben 216fa74d99 misc improvements from blitzpool 2026-05-06 19:05:10 -04:00
Ben 355437df3c test coverage 2026-05-06 18:28:05 -04:00
27 changed files with 9335 additions and 8505 deletions
+1
View File
@@ -0,0 +1 @@
omit=optional
+7539 -8197
View File
File diff suppressed because it is too large Load Diff
+42 -43
View File
@@ -21,62 +21,61 @@
"test:e2e": "jest --config ./test/jest-e2e.json"
},
"dependencies": {
"@nestjs/axios": "^3.0.0",
"@nestjs/cache-manager": "^2.0.1",
"@nestjs/common": "^9.0.0",
"@nestjs/config": "^2.3.2",
"@nestjs/core": "^9.0.0",
"@nestjs/platform-fastify": "^9.4.2",
"@nestjs/schedule": "^3.0.1",
"@nestjs/typeorm": "^10.0.0",
"axios": "^1.4.0",
"big.js": "^6.2.1",
"bitcoin-address-validation": "^2.2.1",
"bitcoinjs-lib": "^6.1.3",
"bitcoinjs-message": "^2.2.0",
"@nestjs/axios": "^4.0.1",
"@nestjs/cache-manager": "^3.1.2",
"@nestjs/common": "^11.1.19",
"@nestjs/config": "^4.0.4",
"@nestjs/core": "^11.1.19",
"@nestjs/platform-fastify": "^11.1.19",
"@nestjs/schedule": "^6.1.3",
"@nestjs/typeorm": "^11.0.1",
"axios": "^1.16.0",
"better-sqlite3": "^12.9.0",
"bitcoin-address-validation": "^2.2.3",
"bitcoinjs-lib": "^6.1.7",
"bs58": "^5.0.0",
"cache-manager": "^5.2.3",
"cache-manager": "^7.2.8",
"class-transformer": "^0.5.1",
"class-validator": "^0.14.0",
"discord.js": "^14.11.0",
"class-validator": "^0.14.4",
"discord.js": "^14.26.4",
"merkle-lib": "^2.0.10",
"node-telegram-bot-api": "^0.61.0",
"pg": "^8.11.3",
"pg": "^8.20.0",
"pg-pubsub": "^0.8.1",
"reflect-metadata": "^0.1.13",
"rpc-bitcoin": "^2.0.0",
"rxjs": "^7.2.0",
"sqlite3": "^5.1.6",
"tiny-secp256k1": "^2.2.3",
"typeorm": "^0.3.17",
"reflect-metadata": "^0.2.2",
"rxjs": "^7.8.2",
"tiny-secp256k1": "^2.2.4",
"typeorm": "^0.3.28",
"zeromq": "6.0.0-beta.19"
},
"devDependencies": {
"@nestjs/cli": "^9.0.0",
"@nestjs/schematics": "^9.0.0",
"@nestjs/testing": "^9.0.0",
"@types/big.js": "^6.1.6",
"@nestjs/cli": "^11.0.21",
"@nestjs/schematics": "^11.1.0",
"@nestjs/testing": "^11.1.19",
"@types/cron": "^2.0.1",
"@types/express": "^4.17.13",
"@types/jest": "29.5.1",
"@types/node": "^18.16.12",
"@types/node-telegram-bot-api": "^0.61.6",
"@types/supertest": "^2.0.11",
"@types/express": "^4.17.25",
"@types/jest": "^29.5.14",
"@types/node": "^18.19.130",
"@types/supertest": "^7.2.0",
"@typescript-eslint/eslint-plugin": "^5.0.0",
"@typescript-eslint/parser": "^5.0.0",
"eslint": "^8.0.1",
"eslint-config-prettier": "^8.3.0",
"eslint-plugin-prettier": "^4.0.0",
"jest": "29.5.0",
"patch-package": "8.0.0",
"prettier": "^2.3.2",
"eslint-config-prettier": "^10.1.8",
"eslint-plugin-prettier": "^5.5.5",
"jest": "^29.7.0",
"patch-package": "^8.0.1",
"prettier": "^3.8.3",
"source-map-support": "^0.5.20",
"supertest": "^6.1.3",
"ts-jest": "29.1.0",
"ts-loader": "^9.2.3",
"ts-node": "^10.0.0",
"supertest": "^7.2.2",
"ts-jest": "^29.4.9",
"ts-loader": "^9.5.7",
"ts-node": "^10.9.2",
"tsconfig-paths": "4.2.0",
"typescript": "^5.0.0"
"typescript": "^5.9.3"
},
"overrides": {
"@fastify/middie": "^9.3.2",
"fastify": "^5.8.5",
"path-to-regexp": "^8.4.2"
},
"jest": {
"moduleFileExtensions": [
-26
View File
@@ -1,26 +0,0 @@
diff --git a/node_modules/rpc-bitcoin/build/src/rpc.d.ts b/node_modules/rpc-bitcoin/build/src/rpc.d.ts
index d25d732..ce4eb22 100644
--- a/node_modules/rpc-bitcoin/build/src/rpc.d.ts
+++ b/node_modules/rpc-bitcoin/build/src/rpc.d.ts
@@ -6,7 +6,7 @@ export declare type RPCIniOptions = RESTIniOptions & {
fullResponse?: boolean;
};
export declare type JSONRPC = {
- jsonrpc?: string | number;
+ jsonrpc?: string;
id?: string | number;
method: string;
params?: object;
diff --git a/node_modules/rpc-bitcoin/build/src/rpc.js b/node_modules/rpc-bitcoin/build/src/rpc.js
index 7fec7aa..3bee5fc 100644
--- a/node_modules/rpc-bitcoin/build/src/rpc.js
+++ b/node_modules/rpc-bitcoin/build/src/rpc.js
@@ -12,7 +12,7 @@ class RPCClient extends rest_1.RESTClient {
}
async rpc(method, params = {}, wallet) {
const uri = typeof wallet === "undefined" ? "/" : "wallet/" + wallet;
- const body = { method, params, jsonrpc: 1.0, id: "rpc-bitcoin" };
+ const body = { method, params, jsonrpc: "1.0", id: "rpc-bitcoin" };
try {
const response = await this.batch(body, uri);
return this.fullResponse ? response : response.result;
@@ -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;
}
@@ -31,6 +34,24 @@ export class AddressSettingsService {
return await this.addressSettingsRepository.update({ address }, { bestDifficulty, bestDifficultyUserAgent });
}
public async updateBestDifficultyIfHigher(address: string, bestDifficulty: number, bestDifficultyUserAgent: string) {
await this.addressSettingsRepository
.createQueryBuilder()
.insert()
.into(AddressSettingsEntity)
.values({ address })
.orIgnore()
.execute();
return await this.addressSettingsRepository
.createQueryBuilder()
.update(AddressSettingsEntity)
.set({ bestDifficulty, bestDifficultyUserAgent })
.where('address = :address', { address })
.andWhere('"bestDifficulty" < :bestDifficulty', { bestDifficulty })
.execute();
}
public async createNew(address: string) {
return await this.addressSettingsRepository.save({ address });
}
@@ -46,10 +67,14 @@ export class AddressSettingsService {
// }
public async resetBestDifficultyAndShares() {
return await this.addressSettingsRepository.update({}, {
shares: 0,
bestDifficulty: 0
});
return await this.addressSettingsRepository
.createQueryBuilder()
.update(AddressSettingsEntity)
.set({
shares: 0,
bestDifficulty: 0
})
.execute();
}
public async getHighScores() {
@@ -259,6 +259,6 @@ export class ClientStatisticsService {
}
public async deleteAll() {
return await this.clientStatisticsRepository.delete({})
return await this.clientStatisticsRepository.clear()
}
}
+15 -1
View File
@@ -140,6 +140,17 @@ export class ClientService {
public async updateBestDifficulty(id: string, bestDifficulty: number) {
return await this.clientRepository.update({ id }, { bestDifficulty });
}
public async updateBestDifficultyIfHigher(id: string, bestDifficulty: number) {
return await this.clientRepository
.createQueryBuilder()
.update(ClientEntity)
.set({ bestDifficulty })
.where('id = :id', { id })
.andWhere('"bestDifficulty" < :bestDifficulty', { bestDifficulty })
.execute();
}
public async connectedClientCount(): Promise<number> {
return await this.clientRepository.count();
}
@@ -173,7 +184,10 @@ export class ClientService {
}
public async deleteAll() {
return await this.clientRepository.softDelete({})
return await this.clientRepository
.createQueryBuilder()
.softDelete()
.execute();
}
// public async getUserAgents() {
+38
View File
@@ -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();
});
});
+106 -7
View File
@@ -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);
}
}
+1 -12
View File
@@ -1,5 +1,4 @@
import { Controller, Patch } from '@nestjs/common';
import * as bitcoinMessage from 'bitcoinjs-message';
@Controller('address')
@@ -7,16 +6,6 @@ export class AddressController {
@Patch('settings')
async settings() {
const publicKey = '...'; // Public key corresponding to the private key used for signing
const message = '...'; // The message that was signed
const signature = '...'; // The signature of the message
const isValid: boolean = bitcoinMessage.verify(message, publicKey, signature);
if (isValid) {
console.log('Signature is valid!');
} else {
console.log('Signature is not valid!');
}
return;
}
}
@@ -13,7 +13,7 @@ describe('ClientController', () => {
const module: TestingModule = await Test.createTestingModule({
imports: [
TypeOrmModule.forRoot({
type: 'sqlite',
type: 'better-sqlite3',
database: ':memory:',
synchronize: true,
autoLoadEntities: true,
+140
View File
@@ -0,0 +1,140 @@
import * as bitcoinjs from 'bitcoinjs-lib';
import { BehaviorSubject, firstValueFrom } from 'rxjs';
import { MockRecording1 } from '../../test/models/MockRecording1';
import { StratumV1JobsService } from '../services/stratum-v1-jobs.service';
import { MiningJob } from './MiningJob';
describe('MiningJob', () => {
let jobTemplate;
let job: MiningJob;
beforeEach(async () => {
jest.useFakeTimers();
jest.setSystemTime(new Date(parseInt(MockRecording1.TIME, 16) * 1000));
const blockTemplate$ = new BehaviorSubject(MockRecording1.BLOCK_TEMPLATE);
const bitcoinRpcService = {
newBlockTemplate$: blockTemplate$.asObservable(),
miningInfo: { blocks: MockRecording1.BLOCK_TEMPLATE.height }
};
jest.spyOn(console, 'log').mockImplementation(() => undefined);
const jobsService = new StratumV1JobsService(bitcoinRpcService as any);
jobTemplate = await firstValueFrom(jobsService.newMiningJob$);
job = new MiningJob(
bitcoinjs.networks.testnet,
'1',
[{ address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', percent: 100 }],
jobTemplate
);
});
afterEach(() => {
jest.restoreAllMocks();
jest.useRealTimers();
});
it('should split coinbase around 12 bytes of extranonce space', () => {
const notify = JSON.parse(job.response(jobTemplate));
const coinbasePart1 = notify.params[2];
const coinbasePart2 = notify.params[3];
const extraNonce1 = '57a6f098';
const extraNonce2 = 'c708000000000000';
const coinbase = bitcoinjs.Transaction.fromHex(`${coinbasePart1}${extraNonce1}${extraNonce2}${coinbasePart2}`);
expect(Buffer.byteLength(extraNonce1 + extraNonce2, 'hex')).toBe(12);
expect(coinbase.ins[0].script.toString('hex')).toContain(`${extraNonce1}${extraNonce2}`);
expect(coinbase.ins[0].script.toString('hex').endsWith(`${extraNonce1}${extraNonce2}`)).toBe(true);
});
it('should update block nonce, timestamp, version mask, and coinbase script', () => {
const extraNonce1 = '57a6f098';
const extraNonce2 = 'c708000000000000';
const timestamp = parseInt(MockRecording1.TIME, 16);
const originalMerkleRoot = Buffer.from(jobTemplate.block.merkleRoot);
const updatedBlock = job.copyAndUpdateBlock(
jobTemplate,
parseInt('00002000', 16),
parseInt('ed460d91', 16),
extraNonce1,
extraNonce2,
timestamp
);
expect(updatedBlock.nonce).toBe(parseInt('ed460d91', 16));
expect(updatedBlock.timestamp).toBe(timestamp);
expect(updatedBlock.version).toBe(jobTemplate.block.version ^ parseInt('00002000', 16));
expect(updatedBlock.transactions[0].ins[0].script.toString('hex').endsWith(`${extraNonce1}${extraNonce2}`)).toBe(true);
expect(updatedBlock.merkleRoot.equals(originalMerkleRoot)).toBe(false);
});
it('should leave block version unchanged without a version mask', () => {
const updatedBlock = job.copyAndUpdateBlock(
jobTemplate,
0,
parseInt('ed460d91', 16),
'57a6f098',
'c708000000000000',
parseInt(MockRecording1.TIME, 16)
);
expect(updatedBlock.version).toBe(jobTemplate.block.version);
});
it('should build the same header as the full block update path', () => {
const versionMask = parseInt('00002000', 16);
const nonce = parseInt('ed460d91', 16);
const extraNonce1 = '57a6f098';
const extraNonce2 = 'c708000000000000';
const timestamp = parseInt(MockRecording1.TIME, 16);
const updatedBlock = job.copyAndUpdateBlock(
jobTemplate,
versionMask,
nonce,
extraNonce1,
extraNonce2,
timestamp
);
const fastHeader = job.buildHeaderBuffer(
jobTemplate,
versionMask,
nonce,
extraNonce1,
extraNonce2,
timestamp
);
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);
});
});
+81 -26
View File
@@ -3,9 +3,10 @@ import * as bitcoinjs from 'bitcoinjs-lib';
import { IJobTemplate } from '../services/stratum-v1-jobs.service';
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;
@@ -16,20 +17,27 @@ export class MiningJob {
private coinbaseTransaction: bitcoinjs.Transaction;
private coinbasePart1: string;
private coinbasePart2: string;
private coinbasePart1Buffer: Buffer;
private coinbasePart2Buffer: Buffer;
private merkleBranchBuffers: Buffer[];
private notifyStaticParams: string;
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();
this.jobTemplateId = jobTemplate.blockData.id;
this.merkleBranchBuffers = jobTemplate.merkle_branch.map(branch => Buffer.from(branch, 'hex'));
this.coinbaseTransaction = this.createCoinbaseTransaction(payoutInformation, jobTemplate.blockData.coinbasevalue);
@@ -41,7 +49,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 +62,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');
@@ -68,8 +87,61 @@ export class MiningJob {
this.coinbasePart1 = serializedCoinbaseTx.slice(0, partOneIndex - (TOTAL_EXTRANONCE_SIZE_BYTES * 2));
this.coinbasePart2 = serializedCoinbaseTx.slice(partOneIndex);
this.coinbasePart1Buffer = Buffer.from(this.coinbasePart1, '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 {
return bitcoinjs.Transaction.fromBuffer(this.coinbaseTransaction.toBuffer());
}
public buildHeaderBuffer(jobTemplate: IJobTemplate, versionMask: number, nonce: number, extraNonce: string, extraNonce2: string, timestamp: number): Buffer {
const coinbaseBuffer = Buffer.concat([
this.coinbasePart1Buffer,
Buffer.from(`${extraNonce}${extraNonce2}`, 'hex'),
this.coinbasePart2Buffer,
]);
const coinbaseHash = bitcoinjs.crypto.hash256(coinbaseBuffer);
const merkleRoot = this.calculateMerkleRootHash(coinbaseHash, this.merkleBranchBuffers);
let version = jobTemplate.block.version;
if (versionMask !== undefined && versionMask != 0) {
version = version ^ versionMask;
}
const header = Buffer.alloc(80);
header.writeInt32LE(version, 0);
jobTemplate.block.prevHash.copy(header, 4);
merkleRoot.copy(header, 36);
header.writeUInt32LE(timestamp, 68);
header.writeUInt32LE(jobTemplate.block.bits, 72);
header.writeUInt32LE(nonce, 76);
return header;
}
public copyAndUpdateBlock(jobTemplate: IJobTemplate, versionMask: number, nonce: number, extraNonce: string, extraNonce2: string, timestamp: number): bitcoinjs.Block {
@@ -79,7 +151,7 @@ export class MiningJob {
return Object.assign(new bitcoinjs.Transaction(), tx);
});
testBlock.transactions[0] = this.coinbaseTransaction;
testBlock.transactions[0] = this.cloneCoinbaseTransaction();
testBlock.nonce = nonce;
@@ -94,7 +166,7 @@ export class MiningJob {
testBlock.transactions[0].ins[0].script = Buffer.from(`${nonceScript.substring(0, nonceScript.length - (TOTAL_EXTRANONCE_SIZE_BYTES * 2))}${extraNonce}${extraNonce2}`, 'hex');
//recompute the root since we updated the coinbase script with the nonces
testBlock.merkleRoot = this.calculateMerkleRootHash(testBlock.transactions[0].getHash(false), jobTemplate.merkle_branch);
testBlock.merkleRoot = this.calculateMerkleRootHash(testBlock.transactions[0].getHash(false), this.merkleBranchBuffers);
testBlock.timestamp = timestamp;
@@ -103,14 +175,14 @@ export class MiningJob {
}
private calculateMerkleRootHash(newRoot: Buffer, merkleBranches: string[]): Buffer {
private calculateMerkleRootHash(newRoot: Buffer, merkleBranches: Buffer[]): Buffer {
const bothMerkles = Buffer.alloc(64);
bothMerkles.set(newRoot);
for (let i = 0; i < merkleBranches.length; i++) {
bothMerkles.set(Buffer.from(merkleBranches[i], 'hex'), 32);
bothMerkles.set(merkleBranches[i], 32);
newRoot = bitcoinjs.crypto.hash256(bothMerkles);
bothMerkles.set(newRoot);
}
@@ -174,24 +246,7 @@ export class MiningJob {
}
public response(jobTemplate: IJobTemplate): string {
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';
return `{"id":null,"method":"${eResponseMethod.MINING_NOTIFY}","params":[${JSON.stringify(this.jobId)},${this.notifyStaticParams}}\n`;
}
+397 -5
View File
@@ -8,6 +8,7 @@ import { DataSource } from 'typeorm';
import { MockRecording1 } from '../../test/models/MockRecording1';
import { AddressSettingsModule } from '../ORM/address-settings/address-settings.module';
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
import { BlocksEntity } from '../ORM/blocks/blocks.entity';
import { BlocksService } from '../ORM/blocks/blocks.service';
import { ClientStatisticsEntity } from '../ORM/client-statistics/client-statistics.entity';
import { ClientStatisticsModule } from '../ORM/client-statistics/client-statistics.module';
@@ -19,7 +20,9 @@ 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 { MiningJob } from './MiningJob';
import { effectiveJobDifficulty, StratumV1Client } from './StratumV1Client';
import { MiningSubmitMessage } from './stratum-messages/MiningSubmitMessage';
@@ -51,6 +54,9 @@ describe('StratumV1Client', () => {
let socketEmitter: (...args: any[]) => void;
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);
@@ -60,12 +66,13 @@ describe('StratumV1Client', () => {
moduleRef = await Test.createTestingModule({
imports: [
TypeOrmModule.forRoot({
type: 'sqlite',
type: 'better-sqlite3',
database: ':memory:',
synchronize: true,
autoLoadEntities: true,
cache: true,
logging: false
logging: false,
entities: [ClientEntity, ClientStatisticsEntity, BlocksEntity]
}),
ClientModule,
ClientStatisticsModule,
@@ -97,18 +104,33 @@ describe('StratumV1Client', () => {
jest.useFakeTimers({ advanceTimers: true })
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);
const dataSource = moduleRef.get<DataSource>(DataSource);
await dataSource.getRepository(ClientStatisticsEntity).delete({});
await dataSource.getRepository(ClientEntity).delete({});
await dataSource.getRepository(ClientStatisticsEntity).clear();
await dataSource.getRepository(ClientEntity).clear();
await dataSource.getRepository(BlocksEntity).clear();
clientStatisticsService = moduleRef.get<ClientStatisticsService>(ClientStatisticsService);
configService = moduleRef.get<ConfigService>(ConfigService);
(configService.get as jest.Mock).mockImplementation((key: string) => {
switch (key) {
case 'DEV_FEE_ADDRESS':
return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
case 'NETWORK':
return 'testnet';
}
return null;
});
(StratumV1Client as any).blockedUserAgentLogState.clear();
(StratumV1Client as any).validationErrorLogState.clear();
bitcoinRpcService = {
newBlockTemplate$: newBlockEmitter.asObservable(),
@@ -130,8 +152,15 @@ describe('StratumV1Client', () => {
});
socket.end = jest.fn();
jest.spyOn(socket, 'destroy').mockImplementation(() => socket);
const addressSettings = moduleRef.get<AddressSettingsService>(AddressSettingsService);
notificationService = {
notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined)
} as any;
blocksService = {
save: jest.fn().mockResolvedValue(undefined)
} as any;
client = new StratumV1Client(
@@ -153,6 +182,9 @@ describe('StratumV1Client', () => {
afterEach(async () => {
client.destroy();
consoleLogSpy.mockRestore();
consoleErrorSpy.mockRestore();
consoleWarnSpy.mockRestore();
jest.useRealTimers();
})
@@ -167,6 +199,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);
@@ -179,6 +232,73 @@ describe('StratumV1Client', () => {
});
it('should block non-compliant user agents on subscribe without allocating a session', async () => {
(configService.get as jest.Mock).mockImplementation((key: string) => {
switch (key) {
case 'NON_COMPLIANT_USER_AGENTS':
return 'NMMiner';
case 'DEV_FEE_ADDRESS':
return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
case 'NETWORK':
return 'testnet';
}
return null;
});
jest.spyOn(socket, 'write').mockImplementation((data) => true);
emitMessage(`{"id":1,"method":"mining.subscribe","params":["NMMiner/1.0"]}`);
await new Promise((r) => setTimeout(r, 1));
expect(socket.destroy).toHaveBeenCalled();
expect(socket.write).not.toHaveBeenCalled();
expect((client as any).statistics).toBeUndefined();
expect(consoleLogSpy).toHaveBeenCalledWith('Blocked non-compliant connection from userAgent: NMMiner');
});
it('should throttle repeated non-compliant user agent logs', async () => {
(configService.get as jest.Mock).mockImplementation((key: string) => {
switch (key) {
case 'NON_COMPLIANT_USER_AGENTS':
return 'NMMiner';
case 'DEV_FEE_ADDRESS':
return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
case 'NETWORK':
return 'testnet';
}
return null;
});
jest.spyOn(socket, 'write').mockImplementation((data) => true);
emitMessage(`{"id":1,"method":"mining.subscribe","params":["NMMiner/1.0"]}`);
await new Promise((r) => setTimeout(r, 1));
const secondSocket = new Socket();
jest.spyOn(secondSocket, 'on').mockImplementation((event: string, listener: (...args: any[]) => void) => {
socketEmitter = listener;
return secondSocket;
});
secondSocket.end = jest.fn();
jest.spyOn(secondSocket, 'destroy').mockImplementation(() => secondSocket);
const secondClient = new StratumV1Client(
secondSocket,
stratumV1JobsService,
bitcoinRpcService,
clientService,
clientStatisticsService,
notificationService,
blocksService,
configService,
moduleRef.get<AddressSettingsService>(AddressSettingsService)
);
socketEmitter(Buffer.from(`{"id":1,"method":"mining.subscribe","params":["NMMiner/1.0"]}\n`));
await new Promise((r) => setTimeout(r, 1));
expect(secondSocket.destroy).toHaveBeenCalled();
expect(consoleLogSpy.mock.calls.filter(call => call[0]?.startsWith('Blocked non-compliant connection'))).toHaveLength(1);
await secondClient.destroy();
});
it('should respond to mining.configure', async () => {
@@ -224,6 +344,7 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
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);
@@ -273,6 +394,277 @@ describe('StratumV1Client', () => {
});
it('should use the header-only fast path for non-block submissions', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
const buildHeaderSpy = jest.spyOn(MiningJob.prototype, 'buildHeaderBuffer');
const fullBlockSpy = jest.spyOn(MiningJob.prototype, 'copyAndUpdateBlock');
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(buildHeaderSpy).toHaveBeenCalled();
expect(fullBlockSpy).not.toHaveBeenCalled();
});
it('should write accepted response before share accounting finishes', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
let finishAccounting: () => void;
const accountingPromise = new Promise<void>((resolve) => {
finishAccounting = resolve;
});
jest.spyOn((client as any).statistics, 'addShares').mockReturnValue(accountingPromise);
emitMessage(MockRecording1.MINING_SUBMIT);
await Promise.resolve();
await Promise.resolve();
await Promise.resolve();
expect((client as any).write).toHaveBeenCalledWith(`{"id":5,"error":null,"result":true}\n`);
finishAccounting();
jest.useRealTimers();
await new Promise((r) => setTimeout(r, 100));
});
it('should update address best difficulty through the atomic path', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
jest.spyOn(client as any, 'calculateDifficulty').mockReturnValue({
submissionDifficulty: 1024,
submissionHash: 'share'
});
const addressSettings = moduleRef.get<AddressSettingsService>(AddressSettingsService);
const getSettingsSpy = jest.spyOn(addressSettings, 'getSettings');
const updateIfHigherSpy = jest.spyOn(addressSettings as any, 'updateBestDifficultyIfHigher').mockResolvedValue({ affected: 1 });
const clientUpdateIfHigherSpy = jest.spyOn(clientService as any, 'updateBestDifficultyIfHigher').mockResolvedValue({ affected: 1 });
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(updateIfHigherSpy).toHaveBeenCalledWith(
'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
1024,
expect.any(String)
);
expect(clientUpdateIfHigherSpy).toHaveBeenCalledWith(expect.any(String), 1024);
expect(getSettingsSpy).not.toHaveBeenCalled();
});
it('should reject duplicate submissions', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
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);
await new Promise((r) => setTimeout(r, 100));
emitMessage(MockRecording1.MINING_SUBMIT);
await new Promise((r) => setTimeout(r, 100));
expect((client as any).write).lastCalledWith(`{"id":5,"result":null,"error":[22,"Duplicate share",""]}\n`);
});
it('should reject submissions for unknown jobs', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
const hashSpy = jest.spyOn(MiningSubmitMessage.prototype, 'hash');
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(`{"id": 5, "method": "mining.submit", "params": ["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.bitaxe3", "ff", "c708000000000000", "64b3f3ec", "ed460d91", "00002000"]}`);
await new Promise((r) => setTimeout(r, 100));
expect((client as any).write).lastCalledWith(`{"id":5,"result":null,"error":[21,"Job not found",""]}\n`);
expect(await clientService.connectedClientCount()).toBe(0);
expect(hashSpy).not.toHaveBeenCalled();
});
it('should reject submissions when the job template has expired', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
stratumV1JobsService.blocks = {};
emitMessage(MockRecording1.MINING_SUBMIT);
await new Promise((r) => setTimeout(r, 100));
expect((client as any).write).lastCalledWith(`{"id":5,"result":null,"error":[21,"Job Template not found",""]}\n`);
});
it('should reject low difficulty shares', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(MockRecording1.MINING_SUGGEST_DIFFICULTY);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
emitMessage(MockRecording1.MINING_SUBMIT);
await new Promise((r) => setTimeout(r, 100));
expect((client as any).write).lastCalledWith(`{"id":5,"result":null,"error":[23,"Difficulty too low",""]}\n`);
expect(await clientService.connectedClientCount()).toBe(0);
});
it('should reject submissions with short extranonce2 values', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
emitMessage(`{"id": 5, "method": "mining.submit", "params": ["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.bitaxe3", "1", "c7080000", "64b3f3ec", "ed460d91", "00002000"]}`);
await new Promise((r) => setTimeout(r, 100));
expect((client as any).write).lastCalledWith(expect.stringContaining(`"error":[20,"Mining Submit validation error"`));
expect(socket.destroy).toHaveBeenCalled();
expect(consoleWarnSpy).toHaveBeenCalledWith(expect.stringContaining('Mining Submit validation error: extraNonce2:isLength'));
});
it('should throttle repeated mining submit validation logs', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
emitMessage(`{"id": 5, "method": "mining.submit", "params": ["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.bitaxe3", "1", "c7080000", "64b3f3ec", "ed460d91", "00002000"]}`);
await new Promise((r) => setTimeout(r, 100));
const secondSocket = new Socket();
jest.spyOn(secondSocket, 'on').mockImplementation((event: string, listener: (...args: any[]) => void) => {
socketEmitter = listener;
return secondSocket;
});
secondSocket.end = jest.fn();
jest.spyOn(secondSocket, 'destroy').mockImplementation(() => secondSocket);
const secondClient = new StratumV1Client(
secondSocket,
stratumV1JobsService,
bitcoinRpcService,
clientService,
clientStatisticsService,
notificationService,
blocksService,
configService,
moduleRef.get<AddressSettingsService>(AddressSettingsService)
);
jest.spyOn(secondClient as any, 'write').mockImplementation((data) => Promise.resolve(true));
jest.spyOn(secondClient as any, 'getRandomHexString').mockReturnValue(MockRecording1.EXTRA_NONCE);
socketEmitter(Buffer.from(`${MockRecording1.MINING_SUBSCRIBE}\n`));
socketEmitter(Buffer.from(`${MockRecording1.MINING_AUTHORIZE}\n`));
await new Promise((r) => setTimeout(r, 100));
socketEmitter(Buffer.from(`{"id": 5, "method": "mining.submit", "params": ["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.bitaxe3", "1", "c7080000", "64b3f3ec", "ed460d91", "00002000"]}\n`));
await new Promise((r) => setTimeout(r, 100));
expect(consoleWarnSpy.mock.calls.filter(call => call[0]?.startsWith('Mining Submit validation error'))).toHaveLength(1);
await secondClient.destroy();
});
it('should close socket when a submit arrives before stratum is initialized', async () => {
const endSpy = jest.spyOn(socket, 'end');
emitMessage(MockRecording1.MINING_SUBMIT);
await new Promise((r) => setTimeout(r, 100));
expect(endSpy).toHaveBeenCalled();
});
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);
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).toHaveBeenCalledWith(expect.objectContaining({
height: MockRecording1.BLOCK_TEMPLATE.height,
minerAddress: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
worker: 'bitaxe3',
sessionId: MockRecording1.EXTRA_NONCE,
blockData: expect.any(String)
}));
expect(notificationService.notifySubscribersBlockFound).toHaveBeenCalled();
expect(addressSettings.resetBestDifficultyAndShares).toHaveBeenCalled();
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);
});
});
+239 -88
View File
@@ -1,8 +1,7 @@
import { ConfigService } from '@nestjs/config';
import Big from 'big.js';
import * as bitcoinjs from 'bitcoinjs-lib';
import { plainToInstance } from 'class-transformer';
import { validate, ValidatorOptions } from 'class-validator';
import { validate, ValidationError, ValidatorOptions } from 'class-validator';
import * as crypto from 'crypto';
import { Socket } from 'net';
import { firstValueFrom, Subscription } from 'rxjs';
@@ -29,20 +28,43 @@ import { EXTRANONCE1_SIZE_BYTES } from './stratum.constants';
import { SuggestDifficulty } from './stratum-messages/SuggestDifficultyMessage';
import { StratumV1ClientStatistics } from './StratumV1ClientStatistics';
const TRUE_DIFF_ONE = 2.695953529101131e67;
const BLOCKED_USER_AGENT_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 {
private static blockedUserAgentLogState = new Map<string, { nextLogAt: number, suppressed: number }>();
private static validationErrorLogState = new Map<string, { nextLogAt: number, suppressed: number, sample: string }>();
public clientSubscription: SubscriptionMessage;
private clientConfiguration: ConfigurationMessage;
private clientAuthorization: AuthorizationMessage;
private clientSuggestedDifficulty: SuggestDifficulty;
private stratumSubscription: Subscription;
private backgroundWork: NodeJS.Timer[] = [];
private backgroundWork: NodeJS.Timeout[] = [];
private statistics: StratumV1ClientStatistics;
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>;
@@ -53,6 +75,7 @@ export class StratumV1Client {
//public hashRate: number = 0;
private buffer: string = '';
private connectionClosed = false;
private miningSubmissionHashes = new Set<string>()
@@ -73,16 +96,19 @@ 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)) {
if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) {
break;
}
try {
await this.handleMessage(m);
} catch (e) {
await this.socket.end();
console.error(e);
}
});
}
})();
});
@@ -139,6 +165,11 @@ export class StratumV1Client {
const errors = await validate(subscriptionMessage, validatorOptions);
if (errors.length === 0) {
if (this.isBlockedUserAgent(subscriptionMessage.userAgent)) {
this.logBlockedUserAgent(subscriptionMessage.userAgent);
this.closeSocket();
return;
}
if (this.sessionStart == null) {
this.sessionStart = new Date();
@@ -224,6 +255,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 +297,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;
@@ -317,17 +350,18 @@ export class StratumV1Client {
} else {
console.log('Mining Submit validation error');
this.logValidationError('Mining Submit validation error', errors);
const err = new StratumErrorMessage(
miningSubmitMessage.id,
eStratumErrorCode.OtherUnknown,
'Mining Submit validation error',
errors).response();
console.error(err);
const success = await this.write(err);
if (!success) {
return;
}
this.closeSocket();
return;
}
break;
}
@@ -352,15 +386,16 @@ export class StratumV1Client {
private async initStratum() {
this.stratumInitialized = true;
if (this.validateHeaderCompliance(this.clientSubscription.userAgent)) {
console.log(`Non compliant connection from userAgent: ${this.clientSubscription.userAgent}`);
await this.socket.end();
if (this.isBlockedUserAgent(this.clientSubscription.userAgent)) {
this.logBlockedUserAgent(this.clientSubscription.userAgent);
this.closeSocket();
return;
}
switch (this.clientSubscription.userAgent) {
case 'cpuminer': {
this.sessionDifficulty = 0.1;
this.oldSessionDifficulty = this.sessionDifficulty;
}
}
@@ -443,7 +478,8 @@ export class StratumV1Client {
network,
this.stratumV1JobsService.getNextId(),
payoutInformation,
jobTemplate
jobTemplate,
this.configService.get('POOL_IDENTIFIER') || 'Public-Pool'
);
this.stratumV1JobsService.addJob(job);
@@ -460,46 +496,28 @@ export class StratumV1Client {
}
private async handleMiningSubmission(submission: MiningSubmitMessage) {
private async ensureClientEntity() {
if (this.clientEntity != null) {
return;
}
if (this.clientEntity == null) {
if (this.creatingEntity == null) {
this.creatingEntity = new Promise(async (resolve, reject) => {
try {
this.clientEntity = await this.clientService.insert({
sessionId: this.extraNonceAndSessionId,
address: this.clientAuthorization.address,
clientName: this.clientAuthorization.worker,
userAgent: this.clientSubscription.userAgent,
startTime: new Date(),
bestDifficulty: 0
});
} catch (e) {
reject(e);
}
resolve();
if (this.creatingEntity == null) {
this.creatingEntity = (async () => {
this.clientEntity = await this.clientService.insert({
sessionId: this.extraNonceAndSessionId,
address: this.clientAuthorization.address,
clientName: this.clientAuthorization.worker,
userAgent: this.clientSubscription.userAgent,
startTime: new Date(),
bestDifficulty: 0
});
await this.creatingEntity;
} else {
await this.creatingEntity;
}
})();
}
const submissionHash = submission.hash();
if(this.miningSubmissionHashes.has(submissionHash)){
const err = new StratumErrorMessage(
submission.id,
eStratumErrorCode.DuplicateShare,
'Duplicate share').response();
const success = await this.write(err);
if (!success) {
return false;
}
return false;
}else{
this.miningSubmissionHashes.add(submissionHash);
}
await this.creatingEntity;
}
private async handleMiningSubmission(submission: MiningSubmitMessage) {
const job = this.stratumV1JobsService.getJobById(submission.jobId);
@@ -532,47 +550,103 @@ 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(
const submissionHash = [
submission.jobId,
submission.extraNonce2,
submission.ntime,
submission.nonce,
submission.versionMask ?? ''
].join(':');
if (this.miningSubmissionHashes.has(submissionHash)) {
const err = new StratumErrorMessage(
submission.id,
eStratumErrorCode.DuplicateShare,
'Duplicate share').response();
const success = await this.write(err);
if (!success) {
return false;
}
return false;
} else {
this.miningSubmissionHashes.add(submissionHash);
}
const versionMask = parseInt(submission.versionMask, 16);
const nonce = parseInt(submission.nonce, 16);
const timestamp = parseInt(submission.ntime, 16);
const header = job.buildHeaderBuffer(
jobTemplate,
parseInt(submission.versionMask, 16),
parseInt(submission.nonce, 16),
versionMask,
nonce,
this.extraNonceAndSessionId,
submission.extraNonce2,
parseInt(submission.ntime, 16)
timestamp
);
const header = updatedJobBlock.toBuffer(true);
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) {
const success = await this.write(JSON.stringify(submission.response()) + '\n');
if (!success) {
return false;
}
if (submissionDifficulty >= jobTemplate.blockData.networkDifficulty) {
console.log('!!! BLOCK FOUND !!!');
const updatedJobBlock = job.copyAndUpdateBlock(
jobTemplate,
versionMask,
nonce,
this.extraNonceAndSessionId,
submission.extraNonce2,
timestamp
);
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}`);
}
}
await this.ensureClientEntity();
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) {
await this.clientService.heartbeatBulkAsync(this.clientEntity.id, this.statistics.hashRate, now);
this.clientService.heartbeatBulkAsync(this.clientEntity.id, this.statistics.hashRate, now);
this.clientEntity.updatedAt = now;
//}
@@ -581,11 +655,9 @@ export class StratumV1Client {
}
if (submissionDifficulty > this.clientEntity.bestDifficulty) {
await this.clientService.updateBestDifficulty(this.clientEntity.id, submissionDifficulty);
await this.clientService.updateBestDifficultyIfHigher(this.clientEntity.id, submissionDifficulty);
this.clientEntity.bestDifficulty = submissionDifficulty;
if (submissionDifficulty > (await this.addressSettingsService.getSettings(this.clientAuthorization.address, true)).bestDifficulty) {
await this.addressSettingsService.updateBestDifficulty(this.clientAuthorization.address, submissionDifficulty, this.clientEntity.userAgent);
}
await this.addressSettingsService.updateBestDifficultyIfHigher(this.clientAuthorization.address, submissionDifficulty, this.clientEntity.userAgent);
}
@@ -605,7 +677,7 @@ export class StratumV1Client {
}
//await this.checkDifficulty();
return true;
return false;
}
@@ -617,6 +689,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 +703,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
@@ -640,34 +720,105 @@ export class StratumV1Client {
const hashResult = bitcoinjs.crypto.hash256(header);
let s64 = this.le256todouble(hashResult);
const truediffone = Big('26959535291011309493156476344723991336010898738574164086137773096960');
const difficulty = truediffone.div(s64.toString());
return { submissionDifficulty: difficulty.toNumber(), submissionHash: hashResult.toString('hex') };
const target = this.le256todouble(hashResult);
const submissionDifficulty = target === 0 ? Number.POSITIVE_INFINITY : TRUE_DIFF_ONE / target;
return { submissionDifficulty, submissionHash: hashResult.toString('hex') };
}
private le256todouble(target: Buffer): bigint {
private le256todouble(target: Buffer): number {
const number = target.reduceRight((acc, byte) => {
// Shift the number 8 bits to the left and OR with the current byte
return (acc << BigInt(8)) | BigInt(byte);
}, BigInt(0));
let number = 0;
for (let i = target.length - 1; i >= 0; i--) {
number = number * 256 + target[i];
}
return number;
}
private validateHeaderCompliance(userAgent: string): boolean {
const headerCompliance = this.configService.get<string>('COMPLIANT_HEADERS');
if (!headerCompliance || headerCompliance.trim() === '') {
private isBlockedUserAgent(userAgent: string): boolean {
const blockedUserAgents = this.configService.get<string>('NON_COMPLIANT_USER_AGENTS')
|| this.configService.get<string>('BLOCKED_USER_AGENTS')
|| this.configService.get<string>('COMPLIANT_HEADERS');
if (!blockedUserAgents || blockedUserAgents.trim() === '') {
return false;
}
const complianceList = headerCompliance.split(',').map(ua => ua.trim().toLowerCase());
const blockedList = blockedUserAgents.split(',').map(ua => ua.trim().toLowerCase());
const userAgentLower = userAgent.toLowerCase();
return complianceList.some(compliant => compliant.length > 0 && userAgentLower.includes(compliant));
return blockedList.some(blocked => blocked.length > 0 && userAgentLower.includes(blocked));
}
private logBlockedUserAgent(userAgent: string) {
const now = Date.now();
const logState = StratumV1Client.blockedUserAgentLogState.get(userAgent);
if (logState != null && now < logState.nextLogAt) {
logState.suppressed += 1;
return;
}
const suppressed = logState?.suppressed ?? 0;
const suffix = suppressed > 0 ? ` (${suppressed} similar connections suppressed)` : '';
console.log(`Blocked non-compliant connection from userAgent: ${userAgent}${suffix}`);
StratumV1Client.blockedUserAgentLogState.set(userAgent, {
nextLogAt: now + BLOCKED_USER_AGENT_LOG_INTERVAL_MS,
suppressed: 0
});
}
private logValidationError(label: string, errors: ValidationError[]) {
const now = Date.now();
const signature = this.getValidationErrorSignature(errors);
const sample = this.getValidationErrorSample(errors);
const key = `${label}:${signature}`;
const logState = StratumV1Client.validationErrorLogState.get(key);
if (logState != null && now < logState.nextLogAt) {
logState.suppressed += 1;
return;
}
const suppressed = logState?.suppressed ?? 0;
const suffix = suppressed > 0 ? ` (${suppressed} similar validation errors suppressed)` : '';
console.warn(`${label}: ${signature}${sample}${suffix}`);
StratumV1Client.validationErrorLogState.set(key, {
nextLogAt: now + VALIDATION_ERROR_LOG_INTERVAL_MS,
suppressed: 0,
sample
});
}
private getValidationErrorSignature(errors: ValidationError[]): string {
if (errors.length === 0) {
return 'unknown';
}
return errors.map(error => {
const constraints = Object.keys(error.constraints ?? {}).sort().join('|') || 'invalid';
return `${error.property}:${constraints}`;
}).join(';');
}
private getValidationErrorSample(errors: ValidationError[]): string {
const values = errors
.map(error => error.value)
.filter(value => value != null)
.map(value => String(value).replace(/[\r\n]/g, '').slice(0, 64));
if (values.length === 0) {
return '';
}
return ` sample=${values.join(',')}`;
}
private closeSocket() {
this.connectionClosed = true;
if (!this.socket.destroyed) {
this.socket.destroy();
}
}
private async write(message: string): Promise<boolean> {
@@ -0,0 +1,107 @@
import { StratumV1ClientStatistics } from './StratumV1ClientStatistics';
describe('StratumV1ClientStatistics', () => {
const client = {
id: 'client-id',
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
clientName: 'bitaxe3',
sessionId: '57a6f098'
} as any;
let clientStatisticsService: {
insert: jest.Mock,
updateBulkAsync: jest.Mock
};
let statistics: StratumV1ClientStatistics;
beforeEach(() => {
jest.useFakeTimers();
jest.setSystemTime(new Date('2026-05-06T12:00:00Z'));
clientStatisticsService = {
insert: jest.fn().mockResolvedValue(undefined),
updateBulkAsync: jest.fn().mockResolvedValue(undefined)
};
statistics = new StratumV1ClientStatistics(clientStatisticsService as any);
});
afterEach(() => {
jest.useRealTimers();
});
it('should insert the first share bucket', async () => {
await statistics.addShares(client, 64);
expect(clientStatisticsService.insert).toHaveBeenCalledWith({
time: new Date('2026-05-06T12:00:00Z').getTime(),
clientId: client.id,
shares: 64,
acceptedCount: 1,
address: client.address,
clientName: client.clientName,
sessionId: client.sessionId
});
});
it('should update the current share bucket for additional shares', async () => {
await statistics.addShares(client, 64);
await statistics.addShares(client, 32);
expect(clientStatisticsService.updateBulkAsync).toHaveBeenCalledWith({
time: new Date('2026-05-06T12:00:00Z').getTime(),
clientId: client.id,
shares: 96,
acceptedCount: 2
});
});
it('should create a new share bucket when the time slot changes', async () => {
await statistics.addShares(client, 64);
jest.setSystemTime(new Date('2026-05-06T12:10:00Z'));
await statistics.addShares(client, 32);
expect(clientStatisticsService.updateBulkAsync).toHaveBeenCalledWith({
time: new Date('2026-05-06T12:00:00Z').getTime(),
clientId: client.id,
shares: 64,
acceptedCount: 1
});
expect(clientStatisticsService.insert).toHaveBeenLastCalledWith({
time: new Date('2026-05-06T12:10:00Z').getTime(),
clientId: client.id,
shares: 32,
acceptedCount: 1,
address: client.address,
clientName: client.clientName,
sessionId: client.sessionId
});
});
it('should not suggest a difficulty change before enough time or shares have passed', () => {
expect(statistics.getSuggestedDifficulty(64)).toBeNull();
});
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(12);
});
it('should increase difficulty for rapid submissions', async () => {
for (let i = 0; i < 5; i++) {
jest.setSystemTime(new Date(Date.parse('2026-05-06T12:00:00Z') + (i * 1000)));
await statistics.addShares(client, 64);
}
expect(statistics.getSuggestedDifficulty(64)).toBe(2048);
});
it('should decrease difficulty for slow submissions', async () => {
for (let i = 0; i < 5; i++) {
jest.setSystemTime(new Date(Date.parse('2026-05-06T12:00:00Z') + (i * 150000)));
await statistics.addShares(client, 64);
}
expect(statistics.getSuggestedDifficulty(128)).toBe(16);
});
});
+24 -18
View File
@@ -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;
}
}
@@ -0,0 +1,40 @@
import { plainToInstance } from 'class-transformer';
import { AuthorizationMessage } from './AuthorizationMessage';
describe('AuthorizationMessage', () => {
it('should parse address, worker, and starting difficulty', () => {
const message = plainToInstance(
AuthorizationMessage,
JSON.parse('{"id":3,"method":"mining.authorize","params":["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.worker1","x,d=2048"]}')
);
expect(message.address).toBe('tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4');
expect(message.worker).toBe('worker1');
expect(message.password).toBe('x,d=2048');
expect(message.startingDiff).toBe(2048);
});
it('should default worker name when one is not provided', () => {
const message = plainToInstance(
AuthorizationMessage,
JSON.parse('{"id":3,"method":"mining.authorize","params":["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4","x"]}')
);
expect(message.worker).toBe('worker');
expect(message.startingDiff).toBeNull();
});
it('should build successful authorization responses', () => {
const message = plainToInstance(
AuthorizationMessage,
JSON.parse('{"id":3,"method":"mining.authorize","params":["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.worker1","x"]}')
);
expect(message.response()).toEqual({
id: 3,
error: null,
result: true
});
});
});
@@ -1,4 +1,5 @@
import { plainToInstance } from 'class-transformer';
import { validate } from 'class-validator';
import { MiningSubmitMessage } from './MiningSubmitMessage';
@@ -29,6 +30,57 @@ describe('MiningSubmitMessage', () => {
expect(message.nonce).toEqual('2402812d');
expect(message.versionMask).toEqual('00006000');
});
it('should validate 8-byte extranonce2 submissions', async () => {
const errors = await validate(message);
expect(errors).toEqual([]);
});
it('should reject short extranonce2 submissions', async () => {
const shortMessage = plainToInstance(
MiningSubmitMessage,
JSON.parse(' {"id": 5, "method": "mining.submit", "params": ["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.bitaxe3", "1", "99020000", "64b1f10f", "2402812d", "00006000"]}'),
);
const errors = await validate(shortMessage);
expect(errors.some(error => error.property === 'extraNonce2')).toBe(true);
});
it('should reject long extranonce2 submissions', async () => {
const longMessage = plainToInstance(
MiningSubmitMessage,
JSON.parse(' {"id": 5, "method": "mining.submit", "params": ["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.bitaxe3", "1", "990200000000000000", "64b1f10f", "2402812d", "00006000"]}'),
);
const errors = await validate(longMessage);
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');
}
@@ -0,0 +1,45 @@
import { plainToInstance } from 'class-transformer';
import { SubscriptionMessage } from './SubscriptionMessage';
describe('SubscriptionMessage', () => {
it('should parse and refine known user agents', () => {
const bosminer = plainToInstance(
SubscriptionMessage,
JSON.parse('{"id":1,"method":"mining.subscribe","params":["bosminer/23.08"]}')
);
const cpuminer = plainToInstance(
SubscriptionMessage,
JSON.parse('{"id":1,"method":"mining.subscribe","params":["cpuminer-opt/1.0"]}')
);
expect(bosminer.userAgent).toBe('Braiins OS');
expect(cpuminer.userAgent).toBe('cpuminer');
});
it('should default missing user agents to unknown', () => {
const message = plainToInstance(
SubscriptionMessage,
JSON.parse('{"id":1,"method":"mining.subscribe","params":[]}')
);
expect(message.userAgent).toBe('unknown');
});
it('should respond with extranonce2 size of 8 bytes', () => {
const message = plainToInstance(
SubscriptionMessage,
JSON.parse('{"id":1,"method":"mining.subscribe","params":["bitaxe v2.2"]}')
);
expect(message.response('57a6f098')).toEqual({
id: 1,
error: null,
result: [
[['mining.notify', '57a6f098']],
'57a6f098',
8
]
});
});
});
+52 -13
View File
@@ -1,6 +1,6 @@
import { Injectable, OnModuleInit } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { RPCClient } from 'rpc-bitcoin';
import axios, { AxiosInstance } from 'axios';
import { asyncScheduler, BehaviorSubject, delay, filter, from, interval, scheduled, shareReplay, startWith, Subject, switchMap } from 'rxjs';
import { RpcBlockService } from '../ORM/rpc-block/rpc-block.service';
import * as zmq from 'zeromq';
@@ -8,15 +8,17 @@ 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 {
private client: RPCClient;
private client: AxiosInstance;
private _newBlockTemplate$: BehaviorSubject<IBlockTemplate> = new BehaviorSubject(undefined);
private pubsubInstance: PGPubsub;
private resetTemplateInterval$ = new Subject<void>();
private rpcRequestId = 0;
public miningInfo: IMiningInfo;
public newBlockTemplate$ = this._newBlockTemplate$.pipe(filter(block => block != null), shareReplay({ refCount: true, bufferSize: 1 }));
@@ -34,14 +36,29 @@ 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');
this.client = new RPCClient({ url, port, timeout, user, pass });
if (cookiefile != null && cookiefile !== '') {
const [cookieUser, cookiePass] = fs.readFileSync(cookiefile).toString().trim().split(':');
user = cookieUser;
pass = cookiePass;
}
this.client.getrpcinfo().then((res) => {
const baseURL = this.buildRpcUrl(url, port);
this.client = axios.create({
baseURL,
timeout,
auth: {
username: user,
password: pass
}
});
this.callRpc('getrpcinfo').then((res) => {
console.log('Bitcoin RPC connected');
}, () => {
console.error('Could not reach RPC host');
@@ -109,13 +126,13 @@ export class BitcoinRpcService implements OnModuleInit {
let blockTemplate: IBlockTemplate;
while (blockTemplate == null) {
blockTemplate = await this.client.getblocktemplate({
template_request: {
blockTemplate = await this.callRpc<IBlockTemplate>('getblocktemplate', [
{
rules: ['segwit'],
mode: 'template',
capabilities: ['serverlist', 'proposal']
}
});
]);
}
try {
@@ -131,7 +148,7 @@ export class BitcoinRpcService implements OnModuleInit {
public async getMiningInfo(): Promise<IMiningInfo> {
try {
return await this.client.getmininginfo();
return await this.callRpc<IMiningInfo>('getmininginfo');
} catch (e) {
console.error('Error getmininginfo', e.message);
return null;
@@ -142,9 +159,7 @@ export class BitcoinRpcService implements OnModuleInit {
public async SUBMIT_BLOCK(hexdata: string): Promise<string> {
let response: string = 'unknown';
try {
response = await this.client.submitblock({
hexdata
});
response = await this.callRpc<string>('submitblock', [hexdata]);
if (response == null) {
response = 'SUCCESS!';
}
@@ -158,4 +173,28 @@ export class BitcoinRpcService implements OnModuleInit {
return response;
}
private async callRpc<T>(method: string, params: unknown[] = []): Promise<T> {
const response = await this.client.post('', {
jsonrpc: '1.0',
id: ++this.rpcRequestId,
method,
params
});
if (response.data.error != null) {
throw response.data.error;
}
return response.data.result;
}
private buildRpcUrl(url: string, port: number): string {
const normalizedUrl = /^https?:\/\//i.test(url) ? url : `http://${url}`;
const rpcUrl = new URL(normalizedUrl);
if (Number.isFinite(port) && port > 0) {
rpcUrl.port = port.toString();
}
return rpcUrl.toString();
}
}
@@ -0,0 +1,112 @@
import { BehaviorSubject, firstValueFrom, skip } from 'rxjs';
import { MockRecording1 } from '../../test/models/MockRecording1';
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
import { StratumV1JobsService } from './stratum-v1-jobs.service';
describe('StratumV1JobsService', () => {
let blockTemplate$: BehaviorSubject<IBlockTemplate>;
let bitcoinRpcService: { newBlockTemplate$: any, miningInfo: { blocks: number } };
let service: StratumV1JobsService;
let consoleLogSpy: jest.SpyInstance;
const createTemplate = (height = MockRecording1.BLOCK_TEMPLATE.height): IBlockTemplate => ({
...MockRecording1.BLOCK_TEMPLATE,
height
});
beforeEach(() => {
jest.useFakeTimers();
jest.setSystemTime(new Date(parseInt(MockRecording1.TIME, 16) * 1000));
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
blockTemplate$ = new BehaviorSubject(createTemplate());
bitcoinRpcService = {
newBlockTemplate$: blockTemplate$.asObservable(),
miningInfo: { blocks: MockRecording1.BLOCK_TEMPLATE.height }
};
service = new StratumV1JobsService(bitcoinRpcService as any);
});
afterEach(() => {
consoleLogSpy.mockRestore();
jest.useRealTimers();
});
it('should create job templates from block templates', async () => {
const jobTemplate = await firstValueFrom(service.newMiningJob$);
expect(jobTemplate.blockData).toEqual(expect.objectContaining({
id: '1',
height: MockRecording1.BLOCK_TEMPLATE.height,
clearJobs: true,
coinbasevalue: MockRecording1.BLOCK_TEMPLATE.coinbasevalue
}));
expect(jobTemplate.merkle_branch.length).toBeGreaterThan(0);
expect(jobTemplate.block.transactions[0].ins[0].witness[0]).toHaveLength(32);
expect(service.getJobTemplateById('1')).toBe(jobTemplate);
});
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;
const nextTemplate = firstValueFrom(service.newMiningJob$.pipe(skip(1)));
blockTemplate$.next(createTemplate(MockRecording1.BLOCK_TEMPLATE.height + 1));
const jobTemplate = await nextTemplate;
expect(jobTemplate.blockData.clearJobs).toBe(true);
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 age retired jobs and templates after the retention window', async () => {
const firstTemplate = await firstValueFrom(service.newMiningJob$);
const oldCreation = Date.now() - (1000 * 60 * 11);
const retiredAt = Date.now() - (1000 * 60 * 11);
firstTemplate.blockData.retiredAt = retiredAt;
for (let i = 4; i >= 0; i--) {
service.jobs[`old-job-${i}`] = { jobId: `old-job-${i}`, creation: oldCreation - i, retiredAt } as any;
(service as any).trackJob(`old-job-${i}`);
service.blocks[`old-template-${i}`] = {
blockData: { creation: oldCreation - i, retiredAt }
} as any;
(service as any).trackBlock(`old-template-${i}`);
}
bitcoinRpcService.miningInfo.blocks = MockRecording1.BLOCK_TEMPLATE.height;
const nextTemplate = firstValueFrom(service.newMiningJob$.pipe(skip(1)));
blockTemplate$.next(createTemplate());
const jobTemplate = await nextTemplate;
expect(jobTemplate.blockData.clearJobs).toBe(false);
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');
service.addJob({ jobId: '1', creation: Date.now() } as any);
expect(service.getNextId()).toBe('2');
expect(service.getJobById('1')).toEqual(expect.objectContaining({ jobId: '1' }));
});
});
+106 -23
View File
@@ -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,11 @@ export class StratumV1JobsService {
public blocks: { [id: number]: IJobTemplate } = {};
private lastBlockHeight = 0;
private jobRetentionMs = parseInt(process.env.JOB_RETENTION_MS) || 600000;
private jobOrder: string[] = [];
private blockOrder: string[] = [];
private jobOrderSet = new Set<string>();
private blockOrderSet = new Set<string>();
constructor(
private readonly bitcoinRpcService: BitcoinRpcService
@@ -109,30 +119,9 @@ 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;
this.trackBlock(data.blockData.id);
}),
shareReplay({ refCount: true, bufferSize: 1 })
)
@@ -162,8 +151,80 @@ export class StratumV1JobsService {
return this.blocks[jobTemplateId];
}
public cleanup(clearJobs: boolean, now: number = Date.now()) {
if (clearJobs) {
for (const id of this.blockOrder) {
const block = this.blocks[id];
if (block != null && block.blockData.retiredAt === undefined) {
block.blockData.retiredAt = now;
}
}
for (const jobId of this.jobOrder) {
const job = this.jobs[jobId];
if (job != null && job.retiredAt === undefined) {
job.retiredAt = now;
}
}
}
this.ageEntries(
this.blocks,
this.blockOrder,
this.blockOrderSet,
now,
entry => entry.blockData.creation,
entry => entry.blockData.retiredAt,
);
this.ageEntries(
this.jobs,
this.jobOrder,
this.jobOrderSet,
now,
entry => entry.creation,
entry => entry.retiredAt,
);
}
private ageEntries<T>(
map: Record<string, T>,
order: string[],
orderSet: Set<string>,
now: number,
getCreation: (entry: T) => number,
getRetiredAt: (entry: T) => number | undefined,
): void {
while (order.length > MIN_RETAINED) {
const id = order[0];
const entry = map[id];
if (entry == null) {
order.shift();
orderSet.delete(id);
continue;
}
const retiredAt = getRetiredAt(entry);
if (retiredAt !== undefined && now - retiredAt > this.jobRetentionMs) {
delete map[id];
order.shift();
orderSet.delete(id);
continue;
}
if (retiredAt === undefined && now - getCreation(entry) > this.jobRetentionMs * 2) {
delete map[id];
order.shift();
orderSet.delete(id);
continue;
}
break;
}
}
public addJob(job: MiningJob) {
this.jobs[job.jobId] = job;
this.trackJob(job.jobId);
this.latestJobId++;
}
@@ -171,6 +232,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);
}
@@ -178,5 +246,20 @@ export class StratumV1JobsService {
return this.latestJobId.toString(16);
}
private trackJob(jobId: string): void {
if (this.jobOrderSet.has(jobId)) {
return;
}
this.jobOrderSet.add(jobId);
this.jobOrder.push(jobId);
}
private trackBlock(blockId: string): void {
if (this.blockOrderSet.has(blockId)) {
return;
}
this.blockOrderSet.add(blockId);
this.blockOrder.push(blockId);
}
}
+78
View File
@@ -0,0 +1,78 @@
import { StratumV1Service } from './stratum-v1.service';
describe('StratumV1Service', () => {
const originalMaster = process.env.MASTER;
const originalStratumPorts = process.env.STRATUM_PORTS;
const originalStratumSecure = process.env.STRATUM_SECURE;
const originalSecureStratumPorts = process.env.SECURE_STRATUM_PORTS;
let service: StratumV1Service;
let clientService;
let consoleLogSpy: jest.SpyInstance;
beforeEach(() => {
jest.useFakeTimers();
clientService = {
deleteAll: jest.fn().mockResolvedValue(undefined)
};
service = new StratumV1Service(
{} as any,
clientService,
{} as any,
{} as any,
{} as any,
{} as any,
{} as any,
{} as any
);
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
});
afterEach(() => {
restoreEnv('MASTER', originalMaster);
restoreEnv('STRATUM_PORTS', originalStratumPorts);
restoreEnv('STRATUM_SECURE', originalStratumSecure);
restoreEnv('SECURE_STRATUM_PORTS', originalSecureStratumPorts);
consoleLogSpy.mockRestore();
jest.useRealTimers();
});
it('should skip Stratum listeners in the master process', async () => {
process.env.MASTER = 'true';
const startSocketServerSpy = jest.spyOn(service as any, 'startSocketServer');
const startSecureSocketServerSpy = jest.spyOn(service as any, 'startSecureSocketServer');
await service.onModuleInit();
jest.runOnlyPendingTimers();
expect(clientService.deleteAll).toHaveBeenCalled();
expect(startSocketServerSpy).not.toHaveBeenCalled();
expect(startSecureSocketServerSpy).not.toHaveBeenCalled();
expect(consoleLogSpy).toHaveBeenCalledWith('Master process skipping Stratum socket listeners');
});
it('should start Stratum listeners in worker processes', async () => {
process.env.MASTER = 'false';
process.env.STRATUM_PORTS = '3333,3334';
process.env.STRATUM_SECURE = 'true';
process.env.SECURE_STRATUM_PORTS = '4333';
const startSocketServerSpy = jest.spyOn(service as any, 'startSocketServer').mockImplementation(() => undefined);
const startSecureSocketServerSpy = jest.spyOn(service as any, 'startSecureSocketServer').mockImplementation(() => undefined);
await service.onModuleInit();
jest.advanceTimersByTime(10000);
expect(clientService.deleteAll).not.toHaveBeenCalled();
expect(startSocketServerSpy).toHaveBeenCalledWith(3333);
expect(startSocketServerSpy).toHaveBeenCalledWith(3334);
expect(startSecureSocketServerSpy).toHaveBeenCalledWith(4333);
});
function restoreEnv(key: string, value: string | undefined) {
if (value == null) {
delete process.env[key];
return;
}
process.env[key] = value;
}
});
+2
View File
@@ -42,6 +42,8 @@ export class StratumV1Service implements OnModuleInit {
if (process.env.MASTER == 'true') {
await this.clientService.deleteAll();
console.log('Master process skipping Stratum socket listeners');
return;
}
// wait for all the other processes to init for an even connection distribution
+63 -22
View File
@@ -1,8 +1,8 @@
import { Injectable, OnModuleInit } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import axios, { AxiosInstance } from 'axios';
import { validate } from 'bitcoin-address-validation';
import { Block } from 'bitcoinjs-lib';
import * as TelegramBot from 'node-telegram-bot-api';
import { TelegramSubscriptionsService } from '../ORM/telegram-subscriptions/telegram-subscriptions.service';
@@ -10,7 +10,9 @@ import { TelegramSubscriptionsService } from '../ORM/telegram-subscriptions/tele
@Injectable()
export class TelegramService implements OnModuleInit {
private bot: TelegramBot;
private bot: AxiosInstance;
private updateOffset = 0;
private pollingTimer: NodeJS.Timeout;
constructor(
private readonly configService: ConfigService,
@@ -20,7 +22,10 @@ export class TelegramService implements OnModuleInit {
if (token == null || token.length < 1) {
return;
}
this.bot = new TelegramBot(token, { polling: true });
this.bot = axios.create({
baseURL: `https://api.telegram.org/bot${token}/`,
timeout: 10000
});
console.log('Telegram bot init');
@@ -32,23 +37,10 @@ export class TelegramService implements OnModuleInit {
return;
}
this.bot.onText(/\/subscribe/, async (msg) => {
const address = msg.text.split('/subscribe ')[1];
if (validate(address) == false) {
this.bot.sendMessage(msg.chat.id, "Invalid address.");
return;
}
await this.telegramSubscriptionsService.saveSubscription(msg.chat.id, address);
this.bot.sendMessage(msg.chat.id, "Subscribed!");
});
this.bot.onText(/\/start/, (msg) => {
this.bot.sendMessage(msg.chat.id, "Welcome to the public-pool bot. /subscribe <address> to get notified.");
});
this.bot.on('message', (msg) => {
console.log(msg);
});
await this.pollUpdates();
this.pollingTimer = setInterval(async () => {
await this.pollUpdates();
}, 2000);
}
public async notifySubscribersBlockFound(address: string, height: number, block: Block, message: string) {
@@ -57,8 +49,57 @@ export class TelegramService implements OnModuleInit {
}
const subscribers = await this.telegramSubscriptionsService.getSubscriptions(address);
subscribers.forEach(subscriber => {
this.bot.sendMessage(subscriber.telegramChatId, `Block Found! Result: ${message}, Height: ${height}`);
await Promise.all(subscribers.map(subscriber => {
return this.sendMessage(subscriber.telegramChatId, `Block Found! Result: ${message}, Height: ${height}`);
}));
}
private async pollUpdates() {
try {
const response = await this.bot.get('getUpdates', {
params: {
offset: this.updateOffset,
timeout: 0
}
});
for (const update of response.data.result ?? []) {
this.updateOffset = update.update_id + 1;
await this.handleMessage(update.message);
}
} catch (e) {
console.error('Telegram polling failed', e.message);
}
}
private async handleMessage(msg: any) {
if (msg?.text == null) {
return;
}
if (msg.text.startsWith('/subscribe')) {
const address = msg.text.split('/subscribe ')[1];
if (validate(address) == false) {
await this.sendMessage(msg.chat.id, 'Invalid address.');
return;
}
await this.telegramSubscriptionsService.saveSubscription(msg.chat.id, address);
await this.sendMessage(msg.chat.id, 'Subscribed!');
return;
}
if (msg.text.startsWith('/start')) {
await this.sendMessage(msg.chat.id, 'Welcome to the public-pool bot. /subscribe <address> to get notified.');
return;
}
console.log(msg);
}
private async sendMessage(chatId: number | string, text: string) {
await this.bot.post('sendMessage', {
chat_id: chatId,
text
});
}
}