Add SV2 DATUM payout hardening

This commit is contained in:
Ben
2026-06-12 23:54:50 -04:00
parent bcd6bb0aa7
commit a33c530268
59 changed files with 7852 additions and 81 deletions
+63 -11
View File
@@ -1,7 +1,8 @@
import { Injectable, OnModuleInit } from '@nestjs/common';
import { Injectable, OnModuleInit, Optional } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import axios, { AxiosInstance } from 'axios';
import { asyncScheduler, BehaviorSubject, delay, filter, from, interval, scheduled, shareReplay, startWith, Subject, switchMap } from 'rxjs';
import { PayoutSnapshotService } from '../ORM/payout-snapshot/payout-snapshot.service';
import { RpcBlockService } from '../ORM/rpc-block/rpc-block.service';
import * as zmq from 'zeromq';
@@ -24,7 +25,9 @@ export class BitcoinRpcService implements OnModuleInit {
constructor(
private readonly configService: ConfigService,
private rpcBlockService: RpcBlockService,
private readonly redisMessagingService: RedisMessagingService
private readonly redisMessagingService: RedisMessagingService,
@Optional()
private readonly payoutSnapshotService?: PayoutSnapshotService
) {
}
@@ -118,7 +121,17 @@ export class BitcoinRpcService implements OnModuleInit {
}
public async getAndBroadcastLatestTemplate() {
if (this.miningInfo?.blocks == null) {
console.warn('Skipping block template broadcast because mining info is not available');
return;
}
const blockTemplate = await this.loadBlockTemplate(this.miningInfo.blocks);
if (blockTemplate == null) {
console.warn(`Skipping block template broadcast for height ${this.miningInfo.blocks}; block template is not available`);
return;
}
this._newBlockTemplate$.next(blockTemplate);
await this.redisMessagingService.setLatestMiningInfo(this.miningInfo);
await this.redisMessagingService.setBlockTemplate(this.miningInfo.blocks, blockTemplate);
@@ -158,13 +171,32 @@ export class BitcoinRpcService implements OnModuleInit {
let blockTemplate: IBlockTemplate;
while (blockTemplate == null) {
blockTemplate = await this.callRpc<IBlockTemplate>('getblocktemplate', [
{
rules: ['segwit'],
mode: 'template',
capabilities: ['serverlist', 'proposal']
}
]);
try {
blockTemplate = await this.callRpc<IBlockTemplate>('getblocktemplate', [
{
rules: ['segwit'],
mode: 'template',
capabilities: ['serverlist', 'proposal']
}
]);
} catch (e) {
console.warn(`Block template is not available yet: ${e.message ?? e}`);
await new Promise(resolve => setTimeout(resolve, 10_000));
}
}
try {
const payoutSnapshot = await this.payoutSnapshotService?.createSnapshotForTemplate({
blockHeight,
coinbaseValueSats: blockTemplate.coinbasevalue,
networkDifficulty: this.calculateNetworkDifficulty(parseInt(blockTemplate.bits, 16)),
});
if (payoutSnapshot != null) {
blockTemplate.payoutSnapshotId = payoutSnapshot.id;
blockTemplate.payoutOutputs = payoutSnapshot.payoutOutputs;
}
} catch (e) {
console.error('Error creating payout snapshot', e);
}
try {
@@ -199,13 +231,25 @@ export class BitcoinRpcService implements OnModuleInit {
console.log(hexdata);
console.log(JSON.stringify(response));
} catch (e) {
response = e;
console.log(`BLOCK SUBMISSION RESPONSE ERROR: ${e}`);
response = e instanceof Error ? e.message : String(e);
console.log(`BLOCK SUBMISSION RESPONSE ERROR: ${response}`);
}
return response;
}
public async TEST_MEMPOOL_ACCEPT(rawTransactions: Buffer[]): Promise<Array<{
txid?: string;
wtxid?: string;
allowed: boolean;
rejectReason?: string;
rejectDetails?: string;
}>> {
return this.callRpc('testmempoolaccept', [
rawTransactions.map(tx => tx.toString('hex')),
]);
}
private async callRpc<T>(method: string, params: unknown[] = []): Promise<T> {
const response = await this.client.post('', {
jsonrpc: '1.0',
@@ -221,6 +265,14 @@ export class BitcoinRpcService implements OnModuleInit {
return response.data.result;
}
private calculateNetworkDifficulty(nBits: number) {
const mantissa: number = nBits & 0x007fffff;
const exponent: number = (nBits >> 24) & 0xff;
const target: number = mantissa * Math.pow(256, (exponent - 3));
const maxTarget = Math.pow(2, 208) * 65535;
return maxTarget / target;
}
private buildRpcUrl(url: string, port: number): string {
const normalizedUrl = /^https?:\/\//i.test(url) ? url : `http://${url}`;
const rpcUrl = new URL(normalizedUrl);
+94
View File
@@ -0,0 +1,94 @@
import { CustomWorkService } from './custom-work.service';
describe('CustomWorkService', () => {
const service = new CustomWorkService();
it('builds SV2 custom-job coinbase splits around the extranonce', () => {
const split = service.buildSv2CoinbaseSplit({
channelId: 1,
requestId: 2,
token: Buffer.from('01', 'hex'),
version: 0x20000000,
prevHash: Buffer.alloc(32),
minNtime: 1,
nBits: 0x170fffff,
coinbaseTxVersion: 2,
coinbasePrefix: Buffer.from('51', 'hex'),
coinbaseTxInputNSequence: 0xffffffff,
coinbaseTxOutputs: Buffer.from('010000000000000000016a', 'hex'),
coinbaseTxLocktime: 0,
merklePath: [],
}, 12);
expect(split.coinbasePrefix.subarray(0, 4).toString('hex')).toBe('02000000');
expect(split.coinbasePrefix.includes(Buffer.from([13]))).toBe(true);
expect(split.coinbasePrefix.subarray(-1).toString('hex')).toBe('51');
expect(split.coinbaseSuffix.subarray(0, 4).toString('hex')).toBe('ffffffff');
expect(split.coinbaseSuffix.subarray(-4).toString('hex')).toBe('00000000');
});
it('encodes Bitcoin varints at boundary values', () => {
expect(service.encodeBitcoinVarInt(0xfc).toString('hex')).toBe('fc');
expect(service.encodeBitcoinVarInt(0xfd).toString('hex')).toBe('fdfd00');
expect(service.encodeBitcoinVarInt(0xffff).toString('hex')).toBe('fdffff');
expect(service.encodeBitcoinVarInt(0x10000).toString('hex')).toBe('fe00000100');
});
it('validates a trivial custom share against a very easy target', () => {
const result = service.validateShare({
coinbasePrefix: Buffer.from('0200000001000000000000000000000000000000000000000000000000000000000000000000000000ffffffff01', 'hex'),
coinbaseSuffix: Buffer.from('ffffffff010000000000000000016a00000000', 'hex'),
merklePath: [],
extranoncePrefix: Buffer.alloc(4),
extranonce: Buffer.alloc(8),
prevHash: Buffer.alloc(32),
nBits: 0x207fffff,
version: 0x20000000,
ntime: 1,
nonce: 1,
shareDifficulty: 0.00000001,
networkDifficulty: Number.MAX_SAFE_INTEGER,
});
expect(result.header).toHaveLength(80);
expect(result.merkleRoot).toHaveLength(32);
expect(result.submissionDifficulty).toBeGreaterThan(0);
expect(result.accepted).toBe(true);
expect(result.isBlockCandidate).toBe(false);
});
it('patches DATUM coinbase target byte before computing the merkle root', () => {
const baseInput = {
coinbasePrefix: Buffer.from('02000000ff', 'hex'),
coinbaseSuffix: Buffer.from('ffffffff010000000000000000016a00000000', 'hex'),
merklePath: [],
extranoncePrefix: Buffer.alloc(0),
extranonce: Buffer.alloc(12),
prevHash: Buffer.alloc(32),
nBits: 0x207fffff,
version: 0x20000000,
ntime: 1,
nonce: 1,
shareDifficulty: 0.00000001,
networkDifficulty: Number.MAX_SAFE_INTEGER,
};
const placeholder = service.validateShare(baseInput);
const patched = service.validateShare({
...baseInput,
coinbaseTargetByteIndex: 4,
coinbaseTargetByte: 0x12,
});
const manuallyPatched = service.computeMerkleRoot(
Buffer.concat([
Buffer.from('0200000012', 'hex'),
Buffer.alloc(12),
baseInput.coinbaseSuffix,
]),
[],
);
expect(patched.merkleRoot.equals(placeholder.merkleRoot)).toBe(false);
expect(patched.merkleRoot.equals(manuallyPatched)).toBe(true);
});
});
+161
View File
@@ -0,0 +1,161 @@
import { Injectable } from '@nestjs/common';
import { BufferReader } from '../models/sv2/sv2-binary-codec';
import { Sv2SetCustomMiningJob } from '../models/sv2/sv2-jdp-messages';
import { DifficultyUtils } from '../utils/difficulty.utils';
import { hash256 } from '../utils/hash.utils';
export interface CustomWorkCoinbaseSplit {
coinbasePrefix: Buffer;
coinbaseSuffix: Buffer;
}
export interface CustomWorkShareValidationInput {
coinbasePrefix: Buffer;
coinbaseSuffix: Buffer;
coinbaseTargetByteIndex?: number;
coinbaseTargetByte?: number;
merklePath: Buffer[];
extranoncePrefix: Buffer;
extranonce: Buffer;
prevHash: Buffer;
nBits: number;
version: number;
ntime: number;
nonce: number;
shareDifficulty: number;
networkDifficulty: number;
}
export interface CustomWorkShareValidationResult {
header: Buffer;
merkleRoot: Buffer;
hashBuffer: Buffer;
submissionDifficulty: number;
accepted: boolean;
isBlockCandidate: boolean;
}
@Injectable()
export class CustomWorkService {
public buildSv2CoinbaseSplit(job: Sv2SetCustomMiningJob, totalExtranonceSize: number): CustomWorkCoinbaseSplit {
const scriptSigLen = job.coinbasePrefix.length + totalExtranonceSize;
const scriptSigLenVarint = this.encodeBitcoinVarInt(scriptSigLen);
const txVersion = Buffer.alloc(4);
txVersion.writeUInt32LE(job.coinbaseTxVersion >>> 0, 0);
const nullOutpoint = Buffer.alloc(36);
nullOutpoint.writeUInt32LE(0xffffffff, 32);
const sequence = Buffer.alloc(4);
sequence.writeUInt32LE(job.coinbaseTxInputNSequence >>> 0, 0);
const locktime = Buffer.alloc(4);
locktime.writeUInt32LE(job.coinbaseTxLocktime >>> 0, 0);
return {
coinbasePrefix: Buffer.concat([
txVersion,
Buffer.from([0x01]),
nullOutpoint,
scriptSigLenVarint,
job.coinbasePrefix,
]),
coinbaseSuffix: Buffer.concat([
sequence,
job.coinbaseTxOutputs,
locktime,
]),
};
}
public validateShare(input: CustomWorkShareValidationInput): CustomWorkShareValidationResult {
const coinbaseTx = Buffer.concat([
input.coinbasePrefix,
input.extranoncePrefix,
input.extranonce,
input.coinbaseSuffix,
]);
if (input.coinbaseTargetByteIndex != null || input.coinbaseTargetByte != null) {
if (input.coinbaseTargetByteIndex == null || input.coinbaseTargetByte == null) {
throw new RangeError('Both coinbase target byte index and value are required');
}
if (input.coinbaseTargetByteIndex < 0 || input.coinbaseTargetByteIndex >= coinbaseTx.length) {
throw new RangeError(`Coinbase target byte index ${input.coinbaseTargetByteIndex} is outside coinbase length ${coinbaseTx.length}`);
}
coinbaseTx[input.coinbaseTargetByteIndex] = input.coinbaseTargetByte & 0xff;
}
const merkleRoot = this.computeMerkleRoot(coinbaseTx, input.merklePath);
const header = this.buildHeader(input.prevHash, merkleRoot, input.version, input.ntime, input.nBits, input.nonce);
const { submissionDifficulty, hashBuffer } = DifficultyUtils.calculateDifficulty(header);
return {
header,
merkleRoot,
hashBuffer,
submissionDifficulty,
accepted: DifficultyUtils.meetsTarget(hashBuffer, DifficultyUtils.difficultyToTarget(input.shareDifficulty)),
isBlockCandidate: DifficultyUtils.meetsTarget(hashBuffer, DifficultyUtils.difficultyToTarget(input.networkDifficulty)),
};
}
public computeMerkleRoot(coinbaseTx: Buffer, merklePath: Buffer[]): Buffer {
let merkleRoot = hash256(coinbaseTx);
const pair = Buffer.alloc(64);
for (const sibling of merklePath) {
pair.fill(0);
merkleRoot.copy(pair, 0);
sibling.copy(pair, 32);
merkleRoot = hash256(pair);
}
return merkleRoot;
}
public buildHeader(prevHash: Buffer, merkleRoot: Buffer, version: number, timestamp: number, nBits: number, nonce: number): Buffer {
const header = Buffer.alloc(80);
header.writeInt32LE(version, 0);
prevHash.copy(header, 4);
merkleRoot.copy(header, 36);
header.writeUInt32LE(timestamp >>> 0, 68);
header.writeUInt32LE(nBits >>> 0, 72);
header.writeUInt32LE(nonce >>> 0, 76);
return header;
}
public encodeBitcoinVarInt(value: number): Buffer {
if (!Number.isSafeInteger(value) || value < 0) {
throw new RangeError(`Invalid Bitcoin varint value ${value}`);
}
if (value < 0xfd) {
return Buffer.from([value]);
}
if (value <= 0xffff) {
const result = Buffer.alloc(3);
result[0] = 0xfd;
result.writeUInt16LE(value, 1);
return result;
}
if (value <= 0xffffffff) {
const result = Buffer.alloc(5);
result[0] = 0xfe;
result.writeUInt32LE(value, 1);
return result;
}
const result = Buffer.alloc(9);
result[0] = 0xff;
result.writeBigUInt64LE(BigInt(value), 1);
return result;
}
public readDatumNullTerminatedString(reader: BufferReader): string {
const bytes: number[] = [];
while (reader.remaining > 0) {
const byte = reader.readU8();
if (byte === 0) {
break;
}
bytes.push(byte);
}
return Buffer.from(bytes).toString('utf8');
}
}
+156
View File
@@ -0,0 +1,156 @@
import { Socket } from 'net';
import { DatumProtocolCommand } from '../models/datum/datum-codec';
import { DatumService } from './datum.service';
import { TemplateProviderService } from './template-provider.service';
function createService(): DatumService {
const templateProvider = {
validateTransactionData: jest.fn(({ transactionList, expectedCount, maxTotalBytes }) => {
if (expectedCount != null && transactionList.length !== expectedCount) {
return { valid: false, errorCode: 'transaction-count-mismatch' };
}
const totalBytes = transactionList.reduce((sum: number, tx: Buffer) => sum + tx.length, 0);
if (maxTotalBytes != null && totalBytes > maxTotalBytes) {
return { valid: false, errorCode: 'transaction-bytes-exceed-limit' };
}
return { valid: true, totalBytes };
}),
};
return new DatumService(
{} as any,
{} as any,
{} as any,
{} as any,
{} as any,
{} as any,
{} as any,
templateProvider as unknown as TemplateProviderService,
undefined,
);
}
function createSocket(): Socket {
return {
destroyed: false,
writableEnded: false,
write: jest.fn((_data: Buffer, callback: (error?: Error) => void) => callback()),
} as any;
}
function createState(cache: any): any {
return {
sessionId: 'datum-test',
session: {
encryptChannelFrame: jest.fn((_command: DatumProtocolCommand, payload: Buffer) => payload),
},
datumJobs: new Map([[7, cache]]),
};
}
describe('DatumService job validation', () => {
it('requests short transaction IDs once a DATUM job advertises transactions', async () => {
const service = createService() as any;
const socket = createSocket();
const cache: any = {
transactionCount: 2,
coinbasePairs: new Map(),
};
const state = createState(cache);
await service.maybeRequestDatumJobValidation(socket, state, 7, cache);
expect(cache.validationState).toBe('requested-short-txids');
expect((socket.write as jest.Mock).mock.calls[0][0].toString('hex')).toBe('501007');
expect(state.session.encryptChannelFrame).toHaveBeenCalledWith(DatumProtocolCommand.MINING, expect.any(Buffer));
});
it('follows short transaction IDs with a full transaction blob request', async () => {
const service = createService() as any;
const socket = createSocket();
const cache: any = {
transactionCount: 1,
coinbasePairs: new Map(),
};
const state = createState(cache);
const shortIdsResponse = Buffer.concat([
Buffer.from([0x90, 7, 0x01]),
Buffer.from('0100', 'hex'),
Buffer.from('010203040506', 'hex'),
Buffer.alloc(32, 0xaa),
Buffer.from([0xfe]),
]);
await service.handleJobValidationResponse(socket, state, shortIdsResponse);
expect(cache.validationState).toBe('requested-full-transaction-blob');
expect(cache.validationShortTxIds.map((id: Buffer) => id.toString('hex'))).toEqual(['010203040506']);
expect((socket.write as jest.Mock).mock.calls[0][0].toString('hex')).toBe('501207');
});
it('fails DATUM job validation when short transaction ID counts disagree', async () => {
const service = createService() as any;
const socket = createSocket();
const warn = jest.spyOn(console, 'warn').mockImplementation(() => undefined);
const cache: any = {
transactionCount: 2,
coinbasePairs: new Map(),
};
const state = createState(cache);
const shortIdsResponse = Buffer.concat([
Buffer.from([0x90, 7, 0x01]),
Buffer.from('0100', 'hex'),
Buffer.from('010203040506', 'hex'),
Buffer.alloc(32, 0xaa),
Buffer.from([0xfe]),
]);
await service.handleJobValidationResponse(socket, state, shortIdsResponse);
expect(cache.validationState).toBe('failed');
expect(cache.validationError).toBe('short txid count mismatch: advertised 2, received 1');
expect(socket.write).not.toHaveBeenCalled();
warn.mockRestore();
});
it('validates and caches full DATUM transaction blobs', async () => {
const service = createService() as any;
const socket = createSocket();
const tx = Buffer.from(
'0100000001'
+ '0000000000000000000000000000000000000000000000000000000000000000'
+ 'ffffffff00ffffffff'
+ '01000000000000000000'
+ '00000000',
'hex',
);
const txSize = Buffer.alloc(3);
txSize.writeUIntLE(tx.length, 0, 3);
const cache: any = {
transactionCount: 1,
totalSize: tx.length,
merkleBranches: [Buffer.alloc(32, 0xbb)],
coinbasePairs: new Map(),
};
const state = createState(cache);
const fullBlobResponse = Buffer.concat([
Buffer.from([0x92, 7, 0x01]),
Buffer.from('0100', 'hex'),
txSize,
tx,
Buffer.from([0xfe]),
]);
await service.handleJobValidationResponse(socket, state, fullBlobResponse);
expect(cache.validationState).toBe('validated');
expect(cache.validationTransactions).toEqual([tx]);
expect(cache.validationError).toBeUndefined();
expect((service as any).templateProvider.validateTransactionData).toHaveBeenCalledWith({
transactionList: [tx],
expectedCount: 1,
maxTotalBytes: tx.length,
expectedCoinbaseMerklePath: [Buffer.alloc(32, 0xbb)],
});
expect(socket.write).not.toHaveBeenCalled();
});
});
+722
View File
@@ -0,0 +1,722 @@
import { Injectable, OnModuleInit } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { getAddressInfo } from 'bitcoin-address-validation';
import * as bitcoinjs from 'bitcoinjs-lib';
import { firstValueFrom } from 'rxjs';
import { Server, Socket } from 'net';
import { BlocksService } from '../ORM/blocks/blocks.service';
import { ClientEntity } from '../ORM/client/client.entity';
import { ClientService } from '../ORM/client/client.service';
import { ShareAccountingService } from '../ORM/share-accounting/share-accounting.service';
import {
DATUM_INITIAL_HEADER_KEY,
DatumFrameReader,
DatumMiningCommand,
DatumProtocolCommand,
DatumRejectReason,
DatumShareResponseStatus,
DatumPowSubmit,
DatumJobValidationStatus,
DatumJobValidationResponse,
deserializeDatumCoinbaserFetch,
deserializeDatumJobValidationResponse,
deserializeDatumMiningCommand,
deserializeDatumPowSubmit,
encodeDatumFrame,
serializeDatumJobValidationFullTransactionBlobRequest,
serializeDatumJobValidationShortTxIdsRequest,
serializeDatumCoinbaserFetchResponse,
serializeDatumShareResponse,
} from '../models/datum/datum-codec';
import { DatumCryptoSession } from '../models/datum/datum-crypto';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { CustomWorkService } from './custom-work.service';
import { RedisMessagingService } from './redis-messaging.service';
import { StratumV1JobsService } from './stratum-v1-jobs.service';
import { DatumKeyPair } from '../models/datum/datum-crypto';
import { TemplateProviderService } from './template-provider.service';
const DEFAULT_DATUM_SHARE_DIFFICULTY = 1;
const DEFAULT_DATUM_PING_INTERVAL_MS = 30_000;
@Injectable()
export class DatumService implements OnModuleInit {
private readonly servers: Server[] = [];
private identityKeys: DatumKeyPair | null = null;
constructor(
private readonly configService: ConfigService,
private readonly jobsService: StratumV1JobsService,
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly clientService: ClientService,
private readonly blocksService: BlocksService,
private readonly shareAccountingService: ShareAccountingService,
private readonly customWorkService: CustomWorkService,
private readonly templateProvider: TemplateProviderService,
private readonly redisMessagingService?: RedisMessagingService,
) {}
public async onModuleInit(): Promise<void> {
if (process.env.API_ONLY === 'true' || process.env.MASTER === 'true') {
return;
}
const ports = this.getPorts();
if (ports.length === 0) {
return;
}
this.identityKeys = this.getIdentityKeys();
console.log(`DATUM server public key: ${Buffer.concat([this.identityKeys.edPublicKey, this.identityKeys.xPublicKey]).toString('hex')}`);
for (const port of ports) {
await this.startServer(port);
}
}
private async startServer(port: number): Promise<void> {
const server = new Server(socket => {
void this.handleSocket(socket);
});
server.on('error', error => {
console.error(`DATUM server error on port ${port}: ${error.message}`);
});
server.listen(port, () => {
console.log(`DATUM server is listening on port ${port}`);
});
this.servers.push(server);
}
private async handleSocket(socket: Socket): Promise<void> {
socket.setKeepAlive(true, 60_000);
socket.setNoDelay(true);
const session = await DatumCryptoSession.create(this.identityKeys ?? undefined);
const reader = new DatumFrameReader(DATUM_INITIAL_HEADER_KEY);
const state: DatumClientState = {
sessionId: Math.random().toString(16).slice(2, 10),
session,
clientEntity: null,
address: null,
workerName: 'default',
userAgent: 'datum/unknown',
pingTimer: null,
datumJobs: new Map(),
};
console.log(`[DATUM ${state.sessionId}] connection accepted from ${socket.remoteAddress}:${socket.remotePort}`);
const close = () => {
void this.destroyClient(state);
if (!socket.destroyed) {
socket.destroy();
}
};
socket.on('data', data => {
void (async () => {
try {
reader.setHeaderXorKey(session.currentReceiveHeaderKey);
const frames = reader.feed(data);
for (const frame of frames) {
let payload = frame.payload;
if (frame.header.isEncryptedChannel) {
payload = session.decryptChannelPayload(payload);
}
await this.handleFrame(socket, state, frame.header.protoCmd, payload);
reader.setHeaderXorKey(session.currentReceiveHeaderKey);
}
} catch (error) {
console.error(`[DATUM ${state.sessionId}] ${error.message}`);
close();
}
})();
});
socket.on('error', close);
socket.on('close', () => void this.destroyClient(state));
}
private async handleFrame(socket: Socket, state: DatumClientState, protoCmd: number, payload: Buffer): Promise<void> {
if (protoCmd === DatumProtocolCommand.HANDSHAKE_INIT) {
const hello = state.session.openHandshake(payload);
state.userAgent = hello.userAgent || 'datum/unknown';
console.log(`[DATUM ${state.sessionId}] handshake init from ${state.userAgent}`);
await this.writeRaw(socket, state.session.buildHandshakeResponse(hello, 'public-pool DATUM'));
console.log(`[DATUM ${state.sessionId}] handshake response sent`);
await this.sendDatumClientConfigure(socket, state);
console.log(`[DATUM ${state.sessionId}] client configure sent`);
this.startDatumPing(socket, state);
return;
}
if (protoCmd !== DatumProtocolCommand.MINING) {
return;
}
const mining = deserializeDatumMiningCommand(payload);
switch (mining.command) {
case DatumMiningCommand.FETCH_COINBASER:
await this.handleCoinbaserFetch(socket, state, mining.body);
break;
case DatumMiningCommand.SUBMIT_POW:
await this.handlePowSubmit(socket, state, mining.body);
break;
case DatumMiningCommand.JOB_VALIDATION:
await this.handleJobValidationResponse(socket, state, mining.body);
break;
case DatumMiningCommand.TEMPLATE_REFRESH:
break;
default:
console.warn(`[DATUM ${state.sessionId}] Ignoring unsupported mining command 0x${mining.command.toString(16)}`);
break;
}
}
private async sendDatumClientConfigure(socket: Socket, state: DatumClientState): Promise<void> {
const payoutAddress = this.getDatumPoolPayoutAddress();
if (!payoutAddress) {
throw new Error('DATUM_POOL_PAYOUT_ADDRESS or DEV_FEE_ADDRESS must be set before DATUM client configuration can be sent');
}
const payoutScript = bitcoinjs.address.toOutputScript(payoutAddress, this.getNetwork());
if (payoutScript.length > 255) {
throw new Error(`DATUM payout script is too long: ${payoutScript.length}`);
}
const coinbaseTag = Buffer.from('public-pool', 'utf8');
if (coinbaseTag.length > 255) {
throw new Error(`DATUM coinbase tag is too long: ${coinbaseTag.length}`);
}
const minDifficulty = Math.max(1, Math.floor(this.getDatumShareDifficulty()));
const vardiffMin = 1n << BigInt(Math.ceil(Math.log2(minDifficulty)));
const vardiffMinBuffer = Buffer.alloc(8);
vardiffMinBuffer.writeBigUInt64LE(vardiffMin, 0);
const payload = Buffer.concat([
Buffer.from([DatumMiningCommand.CLIENT_CONFIGURE]),
Buffer.from([1, payoutScript.length]),
payoutScript,
uint32le(0x50554250),
Buffer.from([coinbaseTag.length]),
coinbaseTag,
vardiffMinBuffer,
Buffer.from([0, 0xfe]),
]);
await this.writeRaw(socket, state.session.encryptChannelFrame(DatumProtocolCommand.MINING, payload, true));
}
private async handleCoinbaserFetch(socket: Socket, state: DatumClientState, payload: Buffer): Promise<void> {
const fetch = deserializeDatumCoinbaserFetch(payload);
const payoutAddress = this.getDatumPoolPayoutAddress();
if (!payoutAddress) {
throw new Error('DATUM_POOL_PAYOUT_ADDRESS or DEV_FEE_ADDRESS must be set before DATUM coinbaser fetches can be served');
}
const scriptPubKey = bitcoinjs.address.toOutputScript(payoutAddress, this.getNetwork());
const response = serializeDatumCoinbaserFetchResponse(fetch.rewardValue, [{
value: fetch.rewardValue,
scriptPubKey,
}]);
await this.writeRaw(socket, state.session.encryptChannelFrame(DatumProtocolCommand.MINING, response));
}
private async handlePowSubmit(socket: Socket, state: DatumClientState, payload: Buffer): Promise<void> {
const pow = deserializeDatumPowSubmit(payload);
const { address, workerName } = this.parseUserIdentity(pow.username);
if (!this.isValidAddress(address)) {
await this.sendShareResponse(socket, state, DatumShareResponseStatus.REJECTED, DatumRejectReason.BAD_USERNAME, pow.nonce, pow.targetByte, pow.jobId);
return;
}
state.address = address;
state.workerName = workerName;
await this.ensureClientEntity(state);
const datumJob = this.updateDatumJobCache(state, pow);
await this.maybeRequestDatumJobValidation(socket, state, pow.jobId, datumJob);
const coinbase = this.getDatumCoinbase(datumJob, pow);
if (coinbase == null || datumJob.prevBlockHash == null || datumJob.nBits == null || datumJob.merkleBranches == null) {
await this.sendShareResponse(socket, state, DatumShareResponseStatus.REJECTED, DatumRejectReason.COINBASE_MISSING, pow.nonce, pow.targetByte, pow.jobId);
return;
}
if (datumJob.targetByteIndex == null) {
await this.sendShareResponse(socket, state, DatumShareResponseStatus.REJECTED, DatumRejectReason.BAD_TARGET, pow.nonce, pow.targetByte, pow.jobId);
return;
}
const latestTemplate = await firstValueFrom(this.jobsService.newMiningJob$);
const templateValidation = this.templateProvider.validateDatumTemplateFastPath({
prevBlockHash: datumJob.prevBlockHash,
nBits: datumJob.nBits,
height: datumJob.height,
version: pow.version,
coinbaseValue: datumJob.coinbaseValue,
totalWeight: datumJob.totalWeight,
totalSize: datumJob.totalSize,
totalSigops: datumJob.totalSigops,
merkleBranches: datumJob.merkleBranches,
});
if (!templateValidation.valid) {
await this.sendShareResponse(
socket,
state,
DatumShareResponseStatus.REJECTED,
this.mapDatumTemplateRejectReason(templateValidation.errorCode),
pow.nonce,
pow.targetByte,
pow.jobId,
);
return;
}
if (datumJob.validationState === 'failed') {
await this.sendShareResponse(socket, state, DatumShareResponseStatus.REJECTED, DatumRejectReason.OTHER, pow.nonce, pow.targetByte, pow.jobId);
return;
}
const shareDifficulty = this.getDatumShareDifficulty();
const nBits = datumJob.nBits.readUInt32LE(0);
const validation = this.customWorkService.validateShare({
coinbasePrefix: coinbase.coinb1,
coinbaseSuffix: coinbase.coinb2,
coinbaseTargetByteIndex: datumJob.targetByteIndex,
coinbaseTargetByte: pow.targetByte,
merklePath: datumJob.merkleBranches,
extranoncePrefix: Buffer.alloc(0),
extranonce: pow.extranonce,
prevHash: datumJob.prevBlockHash,
nBits,
version: pow.version,
ntime: pow.ntime,
nonce: pow.nonce,
shareDifficulty,
networkDifficulty: latestTemplate.blockData.networkDifficulty,
});
if (!validation.accepted) {
await this.sendShareResponse(socket, state, DatumShareResponseStatus.REJECTED, DatumRejectReason.HIGH_HASH, pow.nonce, pow.targetByte, pow.jobId);
return;
}
const isBlockCandidate = validation.isBlockCandidate || pow.isBlock;
const blockSubmissionResult = validation.isBlockCandidate ? 'datum-gateway-submit-expected' : null;
await this.sendShareResponse(socket, state, DatumShareResponseStatus.ACCEPTED, 0, pow.nonce, pow.targetByte, pow.jobId);
if (isBlockCandidate) {
await this.blocksService.save({
height: datumJob.height ?? latestTemplate.blockData.height,
minerAddress: address,
worker: workerName,
sessionId: state.sessionId,
blockData: validation.header.toString('hex'),
blockSubmissionResult,
payoutSnapshotId: null,
});
}
await this.shareAccountingService.recordAcceptedShare({
protocol: 'datum',
workSource: 'miner_template',
workProtocol: 'datum',
address,
clientName: workerName,
sessionId: state.sessionId,
clientId: state.clientEntity.id,
jobId: pow.jobId.toString(16),
jobTemplateId: latestTemplate.blockData.id,
blockHeight: datumJob.height ?? latestTemplate.blockData.height,
creditedDifficulty: shareDifficulty,
submissionDifficulty: validation.submissionDifficulty,
networkDifficulty: latestTemplate.blockData.networkDifficulty,
nonce: pow.nonce,
ntime: pow.ntime,
version: pow.version,
extraNonce2: pow.extranonce.toString('hex'),
isBlockCandidate,
blockSubmissionResult,
});
await this.redisMessagingService?.setClientPresence({
clientId: state.clientEntity.id,
address,
clientName: workerName,
sessionId: state.sessionId,
userAgent: state.userAgent,
startTime: new Date(state.clientEntity.startTime).toISOString(),
lastSeen: new Date().toISOString(),
hashRate: 0,
bestDifficulty: validation.submissionDifficulty,
});
}
private updateDatumJobCache(state: DatumClientState, pow: DatumPowSubmit): DatumJobCache {
let cache = state.datumJobs.get(pow.jobId);
if (cache == null || pow.prevBlockHash != null) {
cache = {
coinbasePairs: new Map(),
};
state.datumJobs.set(pow.jobId, cache);
}
if (pow.prevBlockHash != null) {
cache.prevBlockHash = pow.prevBlockHash;
}
if (pow.targetByteIndex != null) {
cache.targetByteIndex = pow.targetByteIndex;
}
if (pow.nBits != null) {
cache.nBits = pow.nBits;
}
if (pow.coinbaserId != null) {
cache.coinbaserId = pow.coinbaserId;
}
if (pow.height != null) {
cache.height = pow.height;
}
if (pow.coinbaseValue != null) {
cache.coinbaseValue = pow.coinbaseValue;
}
if (pow.transactionCount != null) {
cache.transactionCount = pow.transactionCount;
}
if (pow.totalWeight != null) {
cache.totalWeight = pow.totalWeight;
}
if (pow.totalSize != null) {
cache.totalSize = pow.totalSize;
}
if (pow.totalSigops != null) {
cache.totalSigops = pow.totalSigops;
}
if (pow.merkleBranches != null) {
cache.merkleBranches = pow.merkleBranches;
}
for (const [coinbaseId, coinbase] of pow.coinbasePairs) {
cache.coinbasePairs.set(coinbaseId, coinbase);
}
if (pow.subsidyOnlyCoinbase != null) {
cache.subsidyOnlyCoinbase = pow.subsidyOnlyCoinbase;
}
return cache;
}
private async maybeRequestDatumJobValidation(
socket: Socket,
state: DatumClientState,
jobId: number,
cache: DatumJobCache,
): Promise<void> {
if (cache.transactionCount == null || cache.validationState != null) {
return;
}
if (cache.transactionCount === 0) {
cache.validationState = 'validated';
cache.validationTransactions = [];
return;
}
cache.validationState = 'requested-short-txids';
cache.validationRequestedAt = new Date();
await this.writeRaw(
socket,
state.session.encryptChannelFrame(
DatumProtocolCommand.MINING,
serializeDatumJobValidationShortTxIdsRequest(jobId),
),
);
}
private async handleJobValidationResponse(socket: Socket, state: DatumClientState, payload: Buffer): Promise<void> {
const response = deserializeDatumJobValidationResponse(payload);
const cache = state.datumJobs.get(response.jobId);
if (cache == null) {
console.warn(`[DATUM ${state.sessionId}] Ignoring validation response for unknown job ${response.jobId}`);
return;
}
if (response.status !== DatumJobValidationStatus.SUCCESS) {
cache.validationState = 'failed';
cache.validationError = `gateway returned status 0x${response.status.toString(16)}`;
console.warn(`[DATUM ${state.sessionId}] Job ${response.jobId} validation failed: ${cache.validationError}`);
return;
}
if (response.kind === 'short-txids') {
if (cache.transactionCount != null && response.transactionCount !== cache.transactionCount) {
cache.validationState = 'failed';
cache.validationError = `short txid count mismatch: advertised ${cache.transactionCount}, received ${response.transactionCount}`;
console.warn(`[DATUM ${state.sessionId}] Job ${response.jobId} validation failed: ${cache.validationError}`);
return;
}
cache.validationShortTxIds = response.shortTxIds;
cache.validationShortTxCrosscheck = response.crosscheck ?? undefined;
cache.validationState = 'requested-full-transaction-blob';
await this.writeRaw(
socket,
state.session.encryptChannelFrame(
DatumProtocolCommand.MINING,
serializeDatumJobValidationFullTransactionBlobRequest(response.jobId),
),
);
return;
}
if (response.kind === 'full-transaction-blob' || response.kind === 'full-transactions') {
const validationError = this.validateDatumTransactionResponse(cache, response);
if (validationError != null) {
cache.validationState = 'failed';
cache.validationError = validationError;
console.warn(`[DATUM ${state.sessionId}] Job ${response.jobId} transaction validation failed: ${validationError}`);
return;
}
cache.validationState = 'validated';
cache.validationError = undefined;
cache.validationTransactions = response.transactions;
cache.validationCompletedAt = new Date();
}
}
private validateDatumTransactionResponse(
cache: DatumJobCache,
response: DatumJobValidationResponse,
): string | null {
if (response.kind === 'short-txids') {
return null;
}
if (cache.transactionCount != null && response.transactionCount !== cache.transactionCount) {
return `transaction count mismatch: advertised ${cache.transactionCount}, received ${response.transactionCount}`;
}
const validation = this.templateProvider.validateTransactionData({
transactionList: response.transactions,
expectedCount: cache.transactionCount,
maxTotalBytes: cache.totalSize,
expectedCoinbaseMerklePath: cache.merkleBranches,
});
if (!validation.valid) {
return validation.errorCode ?? 'invalid-transaction-data';
}
return null;
}
private getDatumCoinbase(cache: DatumJobCache, pow: DatumPowSubmit): { coinb1: Buffer; coinb2: Buffer } | undefined {
if (pow.subsidyOnly) {
return cache.subsidyOnlyCoinbase ?? cache.coinbasePairs.get(pow.coinbaseId);
}
return cache.coinbasePairs.get(pow.coinbaseId) ?? cache.subsidyOnlyCoinbase;
}
private async sendShareResponse(
socket: Socket,
state: DatumClientState,
status: DatumShareResponseStatus,
reasonCode: number,
nonce: number,
targetByte: number,
jobId: number,
): Promise<void> {
const payload = serializeDatumShareResponse({ status, reasonCode, nonce, targetByte, jobId });
await this.writeRaw(socket, state.session.encryptChannelFrame(DatumProtocolCommand.MINING, payload));
}
private async ensureClientEntity(state: DatumClientState): Promise<void> {
if (state.clientEntity != null) {
return;
}
state.clientEntity = await this.clientService.insert({
sessionId: state.sessionId,
address: state.address,
clientName: state.workerName,
userAgent: state.userAgent,
startTime: new Date(),
bestDifficulty: 0,
});
}
private async destroyClient(state: DatumClientState): Promise<void> {
this.stopDatumPing(state);
if (state.clientEntity?.id == null) {
return;
}
await this.redisMessagingService?.removeClientPresence(state.clientEntity.id, state.clientEntity.address);
await this.clientService.delete(state.clientEntity.id);
state.clientEntity = null;
}
private startDatumPing(socket: Socket, state: DatumClientState): void {
this.stopDatumPing(state);
state.pingTimer = setInterval(() => {
if (socket.destroyed || socket.writableEnded) {
this.stopDatumPing(state);
return;
}
void this.writeRaw(
socket,
state.session.encryptChannelFrame(DatumProtocolCommand.PING, Buffer.alloc(0)),
).catch(error => {
console.error(`[DATUM ${state.sessionId}] ping failed: ${error.message}`);
if (!socket.destroyed) {
socket.destroy();
}
});
}, this.getDatumPingIntervalMs());
state.pingTimer.unref?.();
}
private stopDatumPing(state: DatumClientState): void {
if (state.pingTimer != null) {
clearInterval(state.pingTimer);
state.pingTimer = null;
}
}
private async writeRaw(socket: Socket, data: Buffer): Promise<void> {
if (socket.destroyed || socket.writableEnded) {
return;
}
await new Promise<void>((resolve, reject) => {
socket.write(data, error => error ? reject(error) : resolve());
});
}
private parseUserIdentity(userIdentity: string): { address: string; workerName: string } {
const parts = userIdentity.split('.');
const address = parts[0] ?? '';
return {
address: this.normalizeAddress(address),
workerName: parts.length > 1 ? parts.slice(1).join('.') : 'default',
};
}
private normalizeAddress(address: string): string {
if (/^(bc1|tb1|bcrt1)/i.test(address)) {
return address.toLowerCase();
}
return address;
}
private isValidAddress(address: string): boolean {
try {
getAddressInfo(address);
return true;
} catch {
return false;
}
}
private getDatumPoolPayoutAddress(): string {
return this.configService.get<string>('DATUM_POOL_PAYOUT_ADDRESS')
|| this.configService.get<string>('PAYOUT_FEE_ADDRESS')
|| this.configService.get<string>('DEV_FEE_ADDRESS')
|| '';
}
private getDatumShareDifficulty(): number {
const configured = parseFloat(this.configService.get<string>('DATUM_SHARE_DIFFICULTY') ?? '');
return Number.isFinite(configured) && configured > 0 ? configured : DEFAULT_DATUM_SHARE_DIFFICULTY;
}
private mapDatumTemplateRejectReason(errorCode?: string): DatumRejectReason {
switch (errorCode) {
case 'prevhash-mismatch':
case 'height-mismatch':
return DatumRejectReason.STALE_BLOCK;
case 'nbits-mismatch':
case 'weight-limit-exceeded':
case 'size-limit-exceeded':
case 'sigop-limit-exceeded':
case 'coinbase-value-too-high':
return DatumRejectReason.OTHER;
case 'version-mismatch':
return DatumRejectReason.BAD_VERSION;
case 'bad-merkle-branch':
return DatumRejectReason.BAD_MERKLE_COUNT;
default:
return DatumRejectReason.OTHER;
}
}
private getDatumPingIntervalMs(): number {
const configured = parseInt(this.configService.get<string>('DATUM_PING_INTERVAL_MS') ?? '', 10);
return Number.isFinite(configured) && configured > 0 ? configured : DEFAULT_DATUM_PING_INTERVAL_MS;
}
private getIdentityKeys(): DatumKeyPair {
const seedHex = this.configService.get<string>('DATUM_IDENTITY_SEED')?.trim();
if (!seedHex) {
const keys = DatumCryptoSession.generateKeyPair();
console.warn('DATUM_IDENTITY_SEED is not set; generated an ephemeral DATUM server identity key');
return keys;
}
if (!/^[0-9a-fA-F]{64}$/.test(seedHex)) {
throw new Error('DATUM_IDENTITY_SEED must be a 32-byte hex string');
}
return DatumCryptoSession.generateKeyPairFromSeed(Buffer.from(seedHex, 'hex'));
}
private getPorts(): number[] {
const configured = this.configService.get<string>('DATUM_PORTS');
if (!configured?.trim()) {
return [];
}
return Array.from(new Set(configured
.split(',')
.map(port => parseInt(port.trim(), 10))
.filter(port => Number.isInteger(port) && port > 0 && port <= 65535)));
}
private getNetwork(): bitcoinjs.networks.Network {
const networkConfig = this.configService.get('NETWORK');
if (networkConfig === 'mainnet') {
return bitcoinjs.networks.bitcoin;
}
if (networkConfig === 'testnet') {
return bitcoinjs.networks.testnet;
}
if (networkConfig === 'regtest') {
return bitcoinjs.networks.regtest;
}
throw new Error('Invalid network configuration');
}
}
interface DatumClientState {
sessionId: string;
session: DatumCryptoSession;
clientEntity: ClientEntity | null;
address: string | null;
workerName: string;
userAgent: string;
pingTimer: NodeJS.Timeout | null;
datumJobs: Map<number, DatumJobCache>;
}
interface DatumJobCache {
prevBlockHash?: Buffer;
targetByteIndex?: number;
nBits?: Buffer;
coinbaserId?: number;
height?: number;
coinbaseValue?: bigint;
transactionCount?: number;
totalWeight?: number;
totalSize?: number;
totalSigops?: number;
merkleBranches?: Buffer[];
coinbasePairs: Map<number, { coinb1: Buffer; coinb2: Buffer }>;
subsidyOnlyCoinbase?: { coinb1: Buffer; coinb2: Buffer };
validationState?: 'requested-short-txids' | 'requested-full-transaction-blob' | 'validated' | 'failed';
validationRequestedAt?: Date;
validationCompletedAt?: Date;
validationShortTxIds?: Buffer[];
validationShortTxCrosscheck?: Buffer;
validationTransactions?: Buffer[];
validationError?: string;
}
function uint32le(value: number): Buffer {
const result = Buffer.alloc(4);
result.writeUInt32LE(value >>> 0, 0);
return result;
}
@@ -91,6 +91,24 @@ describe('StratumV1JobsService', () => {
expect(service.getJobTemplateById(jobTemplate.blockData.id)).toBe(jobTemplate);
});
it('should emit when payout snapshot metadata changes', async () => {
await firstValueFrom(service.newMiningJob$);
const payoutTemplate = createTemplate();
payoutTemplate.payoutSnapshotId = '42';
payoutTemplate.payoutOutputs = [
{ address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', amountSats: 123456 },
];
const nextTemplate = firstValueFrom(service.newMiningJob$.pipe(skip(1)));
blockTemplate$.next(payoutTemplate);
const jobTemplate = await nextTemplate;
expect(jobTemplate.blockData.clearJobs).toBe(false);
expect(jobTemplate.blockData.payoutSnapshotId).toBe('42');
expect(jobTemplate.blockData.payoutOutputs).toEqual(payoutTemplate.payoutOutputs);
});
it('should age old jobs and templates after five minutes', async () => {
await firstValueFrom(service.newMiningJob$);
const oldCreation = Date.now() - (1000 * 60 * 11);
+29 -4
View File
@@ -4,7 +4,8 @@ import * as merkle from 'merkle-lib';
import * as merkleProof from 'merkle-lib/proof';
import { combineLatest, delay, filter, from, interval, map, Observable, shareReplay, startWith, switchMap, tap } from 'rxjs';
import { MiningJob } from '../models/MiningJob';
import { IBlockTemplateTx } from '../models/bitcoin-rpc/IBlockTemplate';
import { AddressObject, MiningJob } from '../models/MiningJob';
import { hash256 } from '../utils/hash.utils';
import { BitcoinRpcService } from './bitcoin-rpc.service';
@@ -19,6 +20,12 @@ export interface IJobTemplate {
networkDifficulty: number;
height: number;
clearJobs: boolean;
payoutSnapshotId?: string;
payoutOutputs?: AddressObject[];
transactions?: IBlockTemplateTx[];
sigoplimit?: number;
sizelimit?: number;
weightlimit?: number;
};
}
@@ -63,6 +70,12 @@ export class StratumV1JobsService {
timestamp,
blockTemplate.height,
blockTemplate.coinbasevalue,
blockTemplate.payoutSnapshotId ?? '',
...(blockTemplate.payoutOutputs ?? []).map(output => [
output.address,
output.amountSats ?? '',
output.percent ?? ''
].join(':')),
...blockTemplate.transactions.map(tx => tx.hash ?? tx.txid ?? tx.data)
].join('|');
@@ -80,11 +93,17 @@ export class StratumV1JobsService {
timestamp,
networkDifficulty: this.calculateNetworkDifficulty(parseInt(blockTemplate.bits, 16)),
clearJobs,
height: blockTemplate.height
height: blockTemplate.height,
payoutSnapshotId: blockTemplate.payoutSnapshotId,
payoutOutputs: blockTemplate.payoutOutputs,
rawTransactions: blockTemplate.transactions,
sigoplimit: blockTemplate.sigoplimit,
sizelimit: blockTemplate.sizelimit,
weightlimit: blockTemplate.weightlimit,
};
}),
filter(next => next != null),
map(({ version, bits, prevHash, transactions, timestamp, coinbasevalue, networkDifficulty, clearJobs, height }) => {
map(({ version, bits, prevHash, transactions, timestamp, coinbasevalue, networkDifficulty, clearJobs, height, payoutSnapshotId, payoutOutputs, rawTransactions, sigoplimit, sizelimit, weightlimit }) => {
const block = new bitcoinjs.Block();
//create an empty coinbase tx
@@ -122,7 +141,13 @@ export class StratumV1JobsService {
coinbasevalue,
networkDifficulty,
height,
clearJobs
clearJobs,
payoutSnapshotId,
payoutOutputs,
transactions: rawTransactions,
sigoplimit,
sizelimit,
weightlimit,
}
}
}),
+9 -5
View File
@@ -9,6 +9,7 @@ import { UserAgentReportService } from '../ORM/_views/user-agent-report/user-age
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
import { BlocksService } from '../ORM/blocks/blocks.service';
import { ClientService } from '../ORM/client/client.service';
import { PayoutSnapshotService } from '../ORM/payout-snapshot/payout-snapshot.service';
import { ShareAccountingService } from '../ORM/share-accounting/share-accounting.service';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { NotificationService } from './notification.service';
@@ -63,7 +64,8 @@ export class StratumV1Service implements OnModuleInit {
private readonly stratumV2Service: StratumV2Service,
private readonly userAgentReportService: UserAgentReportService,
private readonly shareAccountingService?: ShareAccountingService,
private readonly redisMessagingService?: RedisMessagingService
private readonly redisMessagingService?: RedisMessagingService,
private readonly payoutSnapshotService?: PayoutSnapshotService
) {
}
@@ -186,7 +188,7 @@ export class StratumV1Service implements OnModuleInit {
try {
protocol = this.detectProtocol(firstChunk);
if (protocol === 'v1') {
client = this.createV1Client(socket);
client = this.createV1Client(socket, 'sv1');
socket.emit('data', firstChunk);
return;
}
@@ -218,7 +220,7 @@ export class StratumV1Service implements OnModuleInit {
return server;
}
private createV1Client(socket: Socket): StratumV1Client {
private createV1Client(socket: Socket, accountingProtocol: 'sv1' | 'sv1_tls'): StratumV1Client {
return new StratumV1Client(
socket,
this.stratumV1JobsService,
@@ -229,7 +231,9 @@ export class StratumV1Service implements OnModuleInit {
this.configService,
this.addressSettingsService,
this.shareAccountingService,
this.redisMessagingService
this.redisMessagingService,
this.payoutSnapshotService,
accountingProtocol,
);
}
@@ -260,7 +264,7 @@ export class StratumV1Service implements OnModuleInit {
socket.setTimeout(this.getSocketTimeoutMs());
socket.setKeepAlive(true, this.getTcpKeepAliveInitialDelayMs());
const client = this.createV1Client(socket);
const client = this.createV1Client(socket, 'sv1_tls');
let cleanedUp = false;
const cleanup = async (reason: string) => {
+13
View File
@@ -6,6 +6,7 @@ import { Server, Socket } from 'net';
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
import { BlocksService } from '../ORM/blocks/blocks.service';
import { ClientService } from '../ORM/client/client.service';
import { PayoutSnapshotService } from '../ORM/payout-snapshot/payout-snapshot.service';
import { ShareAccountingService } from '../ORM/share-accounting/share-accounting.service';
import { StratumV2Client } from '../models/StratumV2Client';
import { encodeSv2AuthorityPublicKey } from '../models/sv2/sv2-authority-key';
@@ -18,9 +19,11 @@ import {
xOnlyPubKeyFromPriv,
} from '../models/sv2/sv2-noise';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { CustomWorkService } from './custom-work.service';
import { NotificationService } from './notification.service';
import { RedisMessagingService } from './redis-messaging.service';
import { StratumV1JobsService } from './stratum-v1-jobs.service';
import { Sv2JobDeclarationRegistryService } from './sv2-job-declaration-registry.service';
const DEFAULT_SOCKET_TIMEOUT_MS = 1000 * 60 * 60;
const DEFAULT_TCP_KEEPALIVE_INITIAL_DELAY_MS = 1000 * 60;
@@ -44,8 +47,11 @@ export class StratumV2Service implements OnModuleInit {
private readonly configService: ConfigService,
private readonly stratumV1JobsService: StratumV1JobsService,
private readonly addressSettingsService: AddressSettingsService,
private readonly customWorkService: CustomWorkService,
private readonly jobDeclarationRegistry: Sv2JobDeclarationRegistryService,
private readonly shareAccountingService?: ShareAccountingService,
private readonly redisMessagingService?: RedisMessagingService,
private readonly payoutSnapshotService?: PayoutSnapshotService,
) {}
public async onModuleInit(): Promise<void> {
@@ -91,8 +97,11 @@ export class StratumV2Service implements OnModuleInit {
this.blocksService,
this.configService,
this.addressSettingsService,
this.customWorkService,
this.jobDeclarationRegistry,
this.shareAccountingService,
this.redisMessagingService,
this.payoutSnapshotService,
);
}
@@ -132,6 +141,10 @@ export class StratumV2Service implements OnModuleInit {
return this.extranonceManager.minerExtranonceSize;
}
public getExtendedTotalExtranonceSize(): number {
return this.extranonceManager.totalSize;
}
private async initializeNoiseConfig(): Promise<void> {
const configuredAuthorityKey = this.configService.get<string>('SV2_AUTHORITY_PRIVKEY');
this.authorityKeyConfigured = configuredAuthorityKey?.length === 64;
@@ -0,0 +1,58 @@
import { Sv2JobDeclarationRegistryService } from './sv2-job-declaration-registry.service';
describe('Sv2JobDeclarationRegistryService', () => {
it('tracks allocated and declared custom work tokens', () => {
const registry = new Sv2JobDeclarationRegistryService();
const coinbaseOutputs = Buffer.from('010000000000000000016a', 'hex');
const allocated = registry.allocateToken('bc1ptest.worker', coinbaseOutputs);
expect(registry.hasKnownToken(allocated.token)).toBe(true);
expect(allocated.coinbaseOutputs).toEqual(coinbaseOutputs);
expect(allocated.coinbaseOutputs).not.toBe(coinbaseOutputs);
const declared = registry.declareJob({
requestId: 7,
miningJobToken: allocated.token,
version: 0x20000000,
coinbaseTxPrefix: Buffer.concat([Buffer.from('51', 'hex'), Buffer.from('6a', 'hex')]),
coinbaseTxSuffix: Buffer.from('52', 'hex'),
wtxidList: [],
excessData: Buffer.alloc(0),
});
expect(registry.hasKnownToken(declared.token)).toBe(true);
expect(registry.getDeclaredJob(declared.token)).toBe(declared);
expect(declared.originalToken).toEqual(allocated.token);
expect(declared.userIdentifier).toBe('bc1ptest.worker');
});
it('rejects declared jobs that omit allocated coinbase output scripts', () => {
const registry = new Sv2JobDeclarationRegistryService();
const coinbaseOutputs = Buffer.from('010000000000000000016a', 'hex');
const allocated = registry.allocateToken('bc1ptest.worker', coinbaseOutputs);
expect(() => registry.declareJob({
requestId: 7,
miningJobToken: allocated.token,
version: 0x20000000,
coinbaseTxPrefix: Buffer.from('51', 'hex'),
coinbaseTxSuffix: Buffer.from('52', 'hex'),
wtxidList: [],
excessData: Buffer.alloc(0),
})).toThrow('coinbase-output-mismatch');
});
it('rejects declarations using unknown mining job tokens', () => {
const registry = new Sv2JobDeclarationRegistryService();
expect(() => registry.declareJob({
requestId: 8,
miningJobToken: Buffer.from('00', 'hex'),
version: 0x20000000,
coinbaseTxPrefix: Buffer.alloc(0),
coinbaseTxSuffix: Buffer.alloc(0),
wtxidList: [],
excessData: Buffer.alloc(0),
})).toThrow('invalid-mining-job-token');
});
});
@@ -0,0 +1,161 @@
import { Injectable } from '@nestjs/common';
import * as crypto from 'crypto';
import { Sv2DeclareMiningJob } from '../models/sv2/sv2-jdp-messages';
import { TemplateProviderTemplate } from './template-provider.service';
export interface Sv2AllocatedMiningJobToken {
token: Buffer;
userIdentifier: string;
coinbaseOutputs: Buffer;
createdAt: number;
}
export interface Sv2DeclaredMiningJob {
originalToken: Buffer;
token: Buffer;
userIdentifier: string;
job: Sv2DeclareMiningJob;
createdAt: number;
templateId?: bigint;
template?: TemplateProviderTemplate;
validationMode: 'token_only' | 'full_template';
providedTransactions: Buffer[];
}
@Injectable()
export class Sv2JobDeclarationRegistryService {
private readonly allocatedTokens = new Map<string, Sv2AllocatedMiningJobToken>();
private readonly declaredJobs = new Map<string, Sv2DeclaredMiningJob>();
private readonly tokenTtlMs = this.readPositiveInt('SV2_JDP_TOKEN_TTL_MS', 10 * 60 * 1000);
public allocateToken(userIdentifier: string, coinbaseOutputs: Buffer): Sv2AllocatedMiningJobToken {
this.cleanup();
const token = crypto.randomBytes(16);
const entry = {
token,
userIdentifier,
coinbaseOutputs: Buffer.from(coinbaseOutputs),
createdAt: Date.now(),
};
this.allocatedTokens.set(token.toString('hex'), entry);
return entry;
}
public declareJob(job: Sv2DeclareMiningJob, metadata: {
templateId?: bigint;
template?: TemplateProviderTemplate;
validationMode?: 'token_only' | 'full_template';
providedTransactions?: Buffer[];
} = {}): Sv2DeclaredMiningJob {
this.cleanup();
const originalTokenHex = job.miningJobToken.toString('hex');
const allocated = this.allocatedTokens.get(originalTokenHex);
if (allocated == null) {
throw new Error('invalid-mining-job-token');
}
if (!this.coinbaseIncludesRequiredOutputs(job, allocated.coinbaseOutputs)) {
throw new Error('coinbase-output-mismatch');
}
const token = crypto.randomBytes(16);
const declared = {
originalToken: Buffer.from(job.miningJobToken),
token,
userIdentifier: allocated.userIdentifier,
job,
createdAt: Date.now(),
templateId: metadata.templateId,
template: metadata.template,
validationMode: metadata.validationMode ?? 'token_only',
providedTransactions: (metadata.providedTransactions ?? []).map(tx => Buffer.from(tx)),
};
this.declaredJobs.set(token.toString('hex'), declared);
return declared;
}
public hasKnownToken(token: Buffer): boolean {
this.cleanup();
const tokenHex = token.toString('hex');
return this.allocatedTokens.has(tokenHex) || this.declaredJobs.has(tokenHex);
}
public getDeclaredJob(token: Buffer): Sv2DeclaredMiningJob | undefined {
this.cleanup();
return this.declaredJobs.get(token.toString('hex'));
}
private cleanup(): void {
const cutoff = Date.now() - this.tokenTtlMs;
for (const [token, entry] of this.allocatedTokens.entries()) {
if (entry.createdAt < cutoff) {
this.allocatedTokens.delete(token);
}
}
for (const [token, entry] of this.declaredJobs.entries()) {
if (entry.createdAt < cutoff) {
this.declaredJobs.delete(token);
}
}
}
private coinbaseIncludesRequiredOutputs(job: Sv2DeclareMiningJob, coinbaseOutputs: Buffer): boolean {
if (coinbaseOutputs.length === 0) {
return true;
}
const declaredCoinbase = Buffer.concat([job.coinbaseTxPrefix, job.coinbaseTxSuffix]);
for (const script of this.extractOutputScripts(coinbaseOutputs)) {
if (script.length === 0 || declaredCoinbase.indexOf(script) < 0) {
return false;
}
}
return true;
}
private extractOutputScripts(outputs: Buffer): Buffer[] {
const { value: outputCount, offset } = this.readVarInt(outputs, 0);
const scripts: Buffer[] = [];
let cursor = offset;
for (let i = 0; i < outputCount; i++) {
if (cursor + 8 > outputs.length) {
throw new Error('invalid-coinbase-outputs');
}
cursor += 8;
const scriptLength = this.readVarInt(outputs, cursor);
cursor = scriptLength.offset;
if (cursor + scriptLength.value > outputs.length) {
throw new Error('invalid-coinbase-outputs');
}
scripts.push(Buffer.from(outputs.subarray(cursor, cursor + scriptLength.value)));
cursor += scriptLength.value;
}
return scripts;
}
private readVarInt(buffer: Buffer, offset: number): { value: number; offset: number } {
if (offset >= buffer.length) {
throw new Error('invalid-coinbase-outputs');
}
const first = buffer[offset];
if (first < 0xfd) {
return { value: first, offset: offset + 1 };
}
if (first === 0xfd) {
if (offset + 3 > buffer.length) {
throw new Error('invalid-coinbase-outputs');
}
return { value: buffer.readUInt16LE(offset + 1), offset: offset + 3 };
}
if (first === 0xfe) {
if (offset + 5 > buffer.length) {
throw new Error('invalid-coinbase-outputs');
}
return { value: buffer.readUInt32LE(offset + 1), offset: offset + 5 };
}
throw new Error('unsupported-large-coinbase-output-count');
}
private readPositiveInt(name: string, fallback: number): number {
const parsed = parseInt(process.env[name] ?? '', 10);
return Number.isFinite(parsed) && parsed > 0 ? parsed : fallback;
}
}
@@ -0,0 +1,503 @@
import * as bitcoinjs from 'bitcoinjs-lib';
import { BufferReader } from '../models/sv2/sv2-binary-codec';
import { Sv2JdpSetupFlags, Sv2MsgType, Sv2Protocol } from '../models/sv2/sv2-constants';
import {
deserializeDeclareMiningJobError,
deserializeDeclareMiningJobSuccess,
deserializeProvideMissingTransactions,
serializeDeclareMiningJob,
serializeProvideMissingTransactionsSuccess,
serializePushSolution,
} from '../models/sv2/sv2-jdp-messages';
import {
deserializeSetupConnectionSuccess,
serializeSetupConnection,
} from '../models/sv2/sv2-messages';
import { Sv2JobDeclarationConnection } from './sv2-job-declaration.service';
import { Sv2JobDeclarationRegistryService } from './sv2-job-declaration-registry.service';
describe('Sv2JobDeclarationConnection compliance', () => {
it('does not echo Job Declaration setup flags in SetupConnection.Success', async () => {
const { connection, sentFrames } = createConnection();
await (connection as any).handleSetupConnection(serializeSetupConnection({
protocol: Sv2Protocol.JOB_DECLARATION,
minVersion: 2,
maxVersion: 2,
flags: Sv2JdpSetupFlags.DECLARE_TX_DATA,
endpoint_host: 'localhost',
endpoint_port: 34264,
vendor: 'jd-client',
hardwareVersion: '',
firmwareVersion: '',
deviceId: '',
}));
const successFrame = sentFrames.find(frame => frame.msgType === Sv2MsgType.SETUP_CONNECTION_SUCCESS);
const success = deserializeSetupConnectionSuccess(new BufferReader(successFrame.payload));
expect(success.flags).toBe(0);
expect(success.usedVersion).toBe(2);
});
it('rejects DeclareMiningJob unless DECLARE_TX_DATA was negotiated', async () => {
const { connection, sentFrames, registry } = createConnection();
await (connection as any).handleDeclareMiningJob(serializeDeclareMiningJob({
requestId: 9,
miningJobToken: Buffer.from('aa', 'hex'),
version: 0x20000000,
coinbaseTxPrefix: Buffer.alloc(0),
coinbaseTxSuffix: Buffer.alloc(0),
wtxidList: [],
excessData: Buffer.alloc(0),
}));
const errorFrame = sentFrames.find(frame => frame.msgType === Sv2MsgType.JDP_DECLARE_MINING_JOB_ERROR);
const error = deserializeDeclareMiningJobError(new BufferReader(errorFrame.payload));
expect(error.requestId).toBe(9);
expect(error.errorCode).toBe('declare-tx-data-not-negotiated');
expect(registry.declareJob).not.toHaveBeenCalled();
});
it('accepts DeclareMiningJob when all wtxids are known to the template provider', async () => {
const registry = new Sv2JobDeclarationRegistryService();
const allocated = registry.allocateToken('bc1ptest.worker', Buffer.alloc(0));
const coinbaseTxPrefix = Buffer.from('010203', 'hex');
const validateDeclaredWtxids = jest.fn().mockReturnValue({
valid: true,
template: { templateId: 42n },
});
const { connection, sentFrames } = createConnection({
registry,
templateProvider: {
validateDeclaredWtxids,
},
});
(connection as any).declareTxData = true;
await (connection as any).handleDeclareMiningJob(serializeDeclareMiningJob({
requestId: 10,
miningJobToken: allocated.token,
version: 0x20000000,
coinbaseTxPrefix,
coinbaseTxSuffix: Buffer.alloc(0),
wtxidList: [Buffer.alloc(32, 1)],
excessData: Buffer.alloc(0),
}));
const successFrame = sentFrames.find(frame => frame.msgType === Sv2MsgType.JDP_DECLARE_MINING_JOB_SUCCESS);
const success = deserializeDeclareMiningJobSuccess(new BufferReader(successFrame.payload));
expect(success.requestId).toBe(10);
expect(registry.getDeclaredJob(success.newMiningJobToken)?.validationMode).toBe('full_template');
expect(registry.getDeclaredJob(success.newMiningJobToken)?.template).toEqual({ templateId: 42n });
expect(validateDeclaredWtxids).toHaveBeenCalledWith({
version: 0x20000000,
coinbaseTxPrefix,
wtxidList: [Buffer.alloc(32, 1)],
});
});
it('rejects DeclareMiningJob when the declared coinbase prefix fails BIP34 validation', async () => {
const registry = new Sv2JobDeclarationRegistryService();
const allocated = registry.allocateToken('bc1ptest.worker', Buffer.alloc(0));
const { connection, sentFrames } = createConnection({
registry,
templateProvider: {
validateDeclaredWtxids: jest.fn().mockReturnValue({
valid: false,
errorCode: 'invalid-job-param-value-coinbase_tx_prefix',
template: { templateId: 42n },
}),
},
});
(connection as any).declareTxData = true;
await (connection as any).handleDeclareMiningJob(serializeDeclareMiningJob({
requestId: 12,
miningJobToken: allocated.token,
version: 0x20000000,
coinbaseTxPrefix: Buffer.from('/public-pool-sri-jdc-e2e/', 'utf8'),
coinbaseTxSuffix: Buffer.alloc(0),
wtxidList: [],
excessData: Buffer.alloc(0),
}));
const errorFrame = sentFrames.find(frame => frame.msgType === Sv2MsgType.JDP_DECLARE_MINING_JOB_ERROR);
const error = deserializeDeclareMiningJobError(new BufferReader(errorFrame.payload));
expect(error.requestId).toBe(12);
expect(error.errorCode).toBe('invalid-job-param-value-coinbase_tx_prefix');
expect(registry.getDeclaredJob(allocated.token)).toBeUndefined();
});
it('requests missing transactions and finalizes after valid missing transaction response', async () => {
const registry = new Sv2JobDeclarationRegistryService();
const allocated = registry.allocateToken('bc1ptest.worker', Buffer.alloc(0));
const templateProvider = {
validateDeclaredWtxids: jest.fn().mockReturnValue({
valid: false,
errorCode: 'missing-transactions',
template: { templateId: 42n },
unknownTxPositionList: [1],
}),
validateProvidedTransactions: jest.fn().mockReturnValue({ valid: true }),
};
const { connection, sentFrames } = createConnection({ registry, templateProvider });
(connection as any).declareTxData = true;
await (connection as any).handleDeclareMiningJob(serializeDeclareMiningJob({
requestId: 11,
miningJobToken: allocated.token,
version: 0x20000000,
coinbaseTxPrefix: Buffer.alloc(0),
coinbaseTxSuffix: Buffer.alloc(0),
wtxidList: [Buffer.alloc(32, 1), Buffer.alloc(32, 2)],
excessData: Buffer.alloc(0),
}));
const missingFrame = sentFrames.find(frame => frame.msgType === Sv2MsgType.JDP_PROVIDE_MISSING_TRANSACTIONS);
const missing = deserializeProvideMissingTransactions(new BufferReader(missingFrame.payload));
expect(missing.requestId).toBe(11);
expect(missing.unknownTxPositionList).toEqual([1]);
await (connection as any).handleProvideMissingTransactionsSuccess(serializeProvideMissingTransactionsSuccess({
requestId: 11,
transactionList: [Buffer.from('0100000000', 'hex')],
}));
const successFrame = sentFrames.find(frame => frame.msgType === Sv2MsgType.JDP_DECLARE_MINING_JOB_SUCCESS);
const success = deserializeDeclareMiningJobSuccess(new BufferReader(successFrame.payload));
expect(registry.getDeclaredJob(success.newMiningJobToken)?.providedTransactions).toEqual([Buffer.from('0100000000', 'hex')]);
expect(registry.getDeclaredJob(success.newMiningJobToken)?.template).toEqual({ templateId: 42n });
});
it('submits PushSolution before persisting the found block', async () => {
const block = createMockBlock('deadbeef');
const templateProvider = {
getTemplate: jest.fn().mockReturnValue({
height: 4990255,
jobTemplate: {
blockData: {
payoutSnapshotId: '17',
},
},
}),
buildBlockFromDeclaredJobSolution: jest.fn().mockReturnValue(block),
};
const callOrder: string[] = [];
const bitcoinRpcService = {
SUBMIT_BLOCK: jest.fn().mockImplementation(async () => {
callOrder.push('submit');
return null;
}),
};
const blocksService = {
save: jest.fn().mockImplementation(async () => {
callOrder.push('save');
}),
};
const payoutSnapshotService = {
finalizeSnapshotForBlock: jest.fn().mockResolvedValue({ finalized: false, reason: 'disabled' }),
};
const notificationService = {
notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined),
};
const { connection } = createConnection({
templateProvider,
bitcoinRpcService,
blocksService,
payoutSnapshotService,
notificationService,
});
(connection as any).latestDeclaredJob = {
token: Buffer.from('0123456789abcdef0123456789abcdef', 'hex'),
userIdentifier: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.worker1',
job: { requestId: 1 } as any,
templateId: 42n,
validationMode: 'full_template',
providedTransactions: [],
};
await (connection as any).handlePushSolution(serializePushSolution({
extranonce: Buffer.from('aabbccdd', 'hex'),
prevHash: Buffer.alloc(32, 1),
nonce: 1,
ntime: 2,
nBits: 0x1d00ffff,
version: 0x20000000,
}));
expect(bitcoinRpcService.SUBMIT_BLOCK).toHaveBeenCalledWith('deadbeef');
expect(blocksService.save).toHaveBeenCalledWith(expect.objectContaining({
height: 4990255,
minerAddress: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
worker: 'worker1',
sessionId: '01234567',
blockData: 'deadbeef',
payoutSnapshotId: '17',
}));
expect(callOrder).toEqual(['submit', 'save']);
expect(payoutSnapshotService.finalizeSnapshotForBlock).toHaveBeenCalledWith({
payoutSnapshotId: '17',
blockHeight: 4990255,
blockSubmissionResult: null,
});
expect(notificationService.notifySubscribersBlockFound).toHaveBeenCalledWith(
'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
4990255,
block,
null,
);
});
it('submits PushSolution using the declared template snapshot after template cache cleanup', async () => {
const block = createMockBlock('feedface');
const declaredTemplate = {
height: 4990300,
jobTemplate: {
blockData: {
payoutSnapshotId: '23',
},
},
};
const templateProvider = {
getTemplate: jest.fn().mockReturnValue(undefined),
getLatestTemplate: jest.fn().mockReturnValue(undefined),
buildBlockFromDeclaredJobSolution: jest.fn().mockReturnValue(block),
};
const bitcoinRpcService = {
SUBMIT_BLOCK: jest.fn().mockResolvedValue('SUCCESS!'),
};
const blocksService = {
save: jest.fn().mockResolvedValue(undefined),
};
const payoutSnapshotService = {
finalizeSnapshotForBlock: jest.fn().mockResolvedValue({ finalized: false, reason: 'disabled' }),
};
const notificationService = {
notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined),
};
const { connection } = createConnection({
templateProvider,
bitcoinRpcService,
blocksService,
payoutSnapshotService,
notificationService,
});
(connection as any).latestDeclaredJob = {
token: Buffer.from('0123456789abcdef0123456789abcdef', 'hex'),
userIdentifier: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.worker1',
job: { requestId: 1 } as any,
templateId: 42n,
template: declaredTemplate,
validationMode: 'full_template',
providedTransactions: [],
};
await (connection as any).handlePushSolution(serializePushSolution({
extranonce: Buffer.from('aabbccdd', 'hex'),
prevHash: Buffer.alloc(32, 1),
nonce: 1,
ntime: 2,
nBits: 0x1d00ffff,
version: 0x20000000,
}));
expect(templateProvider.getTemplate).not.toHaveBeenCalled();
expect(templateProvider.buildBlockFromDeclaredJobSolution).toHaveBeenCalledWith(expect.objectContaining({
template: declaredTemplate,
}));
expect(bitcoinRpcService.SUBMIT_BLOCK).toHaveBeenCalledWith('feedface');
expect(blocksService.save).toHaveBeenCalledWith(expect.objectContaining({
height: 4990300,
blockSubmissionResult: 'SUCCESS!',
payoutSnapshotId: '23',
}));
});
it('does not resubmit duplicate PushSolution messages', async () => {
const block = createMockBlock('feedface');
const templateProvider = {
buildBlockFromDeclaredJobSolution: jest.fn().mockReturnValue(block),
};
const bitcoinRpcService = {
SUBMIT_BLOCK: jest.fn().mockResolvedValue('SUCCESS!'),
};
const { connection } = createConnection({
templateProvider,
bitcoinRpcService,
blocksService: { save: jest.fn().mockResolvedValue(undefined) },
});
(connection as any).latestDeclaredJob = {
token: Buffer.from('0123456789abcdef0123456789abcdef', 'hex'),
userIdentifier: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.worker1',
job: { requestId: 1 } as any,
template: {
height: 4990300,
jobTemplate: { blockData: { payoutSnapshotId: null } },
},
validationMode: 'full_template',
providedTransactions: [],
};
const payload = serializePushSolution({
extranonce: Buffer.from('aabbccdd', 'hex'),
prevHash: Buffer.alloc(32, 1),
nonce: 1,
ntime: 2,
nBits: 0x1d00ffff,
version: 0x20000000,
});
await (connection as any).handlePushSolution(payload);
await (connection as any).handlePushSolution(payload);
expect(bitcoinRpcService.SUBMIT_BLOCK).toHaveBeenCalledTimes(1);
});
it('does not submit PushSolution blocks with a coinbase missing the BIP34 height', async () => {
const warn = jest.spyOn(console, 'warn').mockImplementation(() => undefined);
const invalidBlock = createBlockWithCoinbaseScript(Buffer.from('/public-pool-sri-jdc-e2e/', 'utf8'));
const templateProvider = {
buildBlockFromDeclaredJobSolution: jest.fn().mockReturnValue(invalidBlock),
validateCoinbaseTransactionHeight: jest.fn().mockReturnValue({
valid: false,
errorCode: 'invalid-job-param-value-coinbase_tx_prefix',
}),
};
const bitcoinRpcService = {
SUBMIT_BLOCK: jest.fn().mockResolvedValue('SUCCESS!'),
};
const blocksService = {
save: jest.fn().mockResolvedValue(undefined),
};
const { connection } = createConnection({
templateProvider,
bitcoinRpcService,
blocksService,
});
(connection as any).latestDeclaredJob = {
token: Buffer.from('0123456789abcdef0123456789abcdef', 'hex'),
userIdentifier: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.worker1',
job: { requestId: 1 } as any,
template: {
height: 4990300,
jobTemplate: { blockData: { payoutSnapshotId: null } },
},
validationMode: 'full_template',
providedTransactions: [],
};
await (connection as any).handlePushSolution(serializePushSolution({
extranonce: Buffer.from('aabbccdd', 'hex'),
prevHash: Buffer.alloc(32, 1),
nonce: 1,
ntime: 2,
nBits: 0x1d00ffff,
version: 0x20000000,
}));
expect(bitcoinRpcService.SUBMIT_BLOCK).not.toHaveBeenCalled();
expect(blocksService.save).not.toHaveBeenCalled();
expect(warn).toHaveBeenCalledWith(expect.stringContaining('invalid-job-param-value-coinbase_tx_prefix'));
warn.mockRestore();
});
function createConnection(overrides: {
registry?: any;
templateProvider?: any;
bitcoinRpcService?: any;
blocksService?: any;
payoutSnapshotService?: any;
notificationService?: any;
} = {}): {
connection: Sv2JobDeclarationConnection;
sentFrames: any[];
registry: { allocateToken: jest.Mock; declareJob: jest.Mock };
} {
const socket = {
setKeepAlive: jest.fn(),
setNoDelay: jest.fn(),
on: jest.fn(),
write: jest.fn((data, callback) => callback?.()),
destroy: jest.fn(),
destroyed: false,
writableEnded: false,
};
const registry = overrides.registry ?? {
allocateToken: jest.fn(),
declareJob: jest.fn(),
};
const templateProvider = {
validateDeclaredWtxids: jest.fn().mockReturnValue({ valid: false, errorCode: 'template-not-found' }),
validateProvidedTransactions: jest.fn().mockReturnValue({ valid: false, errorCode: 'template-not-found' }),
validateCoinbaseTransactionHeight: jest.fn().mockReturnValue({ valid: true }),
...overrides.templateProvider,
};
const connection = new Sv2JobDeclarationConnection(
socket as any,
{
getNoiseConfig: () => ({
staticKeypair: {
privateKey: Buffer.alloc(32),
publicKey: Buffer.alloc(64),
},
certificateMessage: {
version: 0,
validFrom: 0,
notValidAfter: 0,
signature: Buffer.alloc(64),
},
}),
} as any,
registry as any,
{} as any,
templateProvider,
overrides.bitcoinRpcService ?? {
SUBMIT_BLOCK: jest.fn().mockResolvedValue(null),
},
overrides.blocksService ?? {
save: jest.fn().mockResolvedValue(undefined),
},
overrides.payoutSnapshotService ?? {
finalizeSnapshotForBlock: jest.fn().mockResolvedValue({ finalized: false, reason: 'disabled' }),
},
overrides.notificationService ?? {
notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined),
},
'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
bitcoinjs.networks.testnet,
);
const sentFrames: any[] = [];
(connection as any).sendFrame = jest.fn((msgType: number, payload: Buffer) => {
sentFrames.push({ msgType, payload });
return Promise.resolve();
});
return { connection, sentFrames, registry: registry as any };
}
function createBlockWithCoinbaseScript(script: Buffer): bitcoinjs.Block {
const tx = new bitcoinjs.Transaction();
tx.version = 2;
tx.addInput(Buffer.alloc(32), 0xffffffff, 0xffffffff, script);
tx.addOutput(Buffer.from('6a', 'hex'), 0);
const block = new bitcoinjs.Block();
block.version = 0x20000000;
block.prevHash = Buffer.alloc(32);
block.timestamp = 2;
block.bits = 0x1d00ffff;
block.nonce = 1;
block.transactions = [tx];
block.merkleRoot = bitcoinjs.Block.calculateMerkleRoot(block.transactions, false);
return block;
}
function createMockBlock(hex: string): { toHex: jest.Mock; transactions: Array<{ toBuffer: jest.Mock }> } {
return {
toHex: jest.fn().mockReturnValue(hex),
transactions: [{
toBuffer: jest.fn().mockReturnValue(Buffer.from('0100000000', 'hex')),
}],
};
}
});
+473
View File
@@ -0,0 +1,473 @@
import { Injectable, OnModuleInit } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import * as bitcoinjs from 'bitcoinjs-lib';
import { Server, Socket } from 'net';
import { BlocksService } from '../ORM/blocks/blocks.service';
import { PayoutSnapshotService } from '../ORM/payout-snapshot/payout-snapshot.service';
import { BufferReader } from '../models/sv2/sv2-binary-codec';
import { SV2_NOISE_ACT1_SIZE, Sv2JdpSetupFlags, Sv2MsgType, Sv2Protocol } from '../models/sv2/sv2-constants';
import { Sv2FrameReader, Sv2FrameWriter } from '../models/sv2/sv2-frame';
import {
deserializeAllocateMiningJobToken,
deserializeDeclareMiningJob,
deserializeProvideMissingTransactionsSuccess,
deserializePushSolution,
serializeAllocateMiningJobTokenSuccess,
serializeDeclareMiningJobError,
serializeDeclareMiningJobSuccess,
serializeProvideMissingTransactions,
Sv2DeclareMiningJob,
} from '../models/sv2/sv2-jdp-messages';
import {
deserializeSetupConnection,
serializeSetupConnectionError,
serializeSetupConnectionSuccess,
} from '../models/sv2/sv2-messages';
import { Sv2NoiseSession } from '../models/sv2/sv2-noise';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { CustomWorkService } from './custom-work.service';
import { NotificationService } from './notification.service';
import { StratumV2Service } from './stratum-v2.service';
import { Sv2DeclaredMiningJob, Sv2JobDeclarationRegistryService } from './sv2-job-declaration-registry.service';
import { TemplateProviderService } from './template-provider.service';
@Injectable()
export class Sv2JobDeclarationService implements OnModuleInit {
private readonly servers: Server[] = [];
constructor(
private readonly configService: ConfigService,
private readonly stratumV2Service: StratumV2Service,
private readonly registry: Sv2JobDeclarationRegistryService,
private readonly customWorkService: CustomWorkService,
private readonly templateProvider: TemplateProviderService,
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly blocksService: BlocksService,
private readonly payoutSnapshotService: PayoutSnapshotService,
private readonly notificationService: NotificationService,
) {}
public async onModuleInit(): Promise<void> {
if (process.env.API_ONLY === 'true' || process.env.MASTER === 'true') {
return;
}
const ports = this.getPorts();
if (ports.length === 0) {
return;
}
await this.stratumV2Service.ensureInitialized();
for (const port of ports) {
this.startServer(port);
}
}
private startServer(port: number): void {
const server = new Server(socket => {
void new Sv2JobDeclarationConnection(
socket,
this.stratumV2Service,
this.registry,
this.customWorkService,
this.templateProvider,
this.bitcoinRpcService,
this.blocksService,
this.payoutSnapshotService,
this.notificationService,
this.getPoolPayoutAddress(),
this.getNetwork(),
).start();
});
server.on('error', error => console.error(`SV2 JDP server error on port ${port}: ${error.message}`));
server.listen(port, () => console.log(`SV2 Job Declaration server is listening on port ${port}`));
this.servers.push(server);
}
private getPoolPayoutAddress(): string {
return this.configService.get<string>('SV2_JDP_POOL_PAYOUT_ADDRESS')
|| this.configService.get<string>('DATUM_POOL_PAYOUT_ADDRESS')
|| this.configService.get<string>('PAYOUT_FEE_ADDRESS')
|| this.configService.get<string>('DEV_FEE_ADDRESS')
|| '';
}
private getPorts(): number[] {
const configured = this.configService.get<string>('SV2_JDP_PORTS');
if (!configured?.trim()) {
return [];
}
return Array.from(new Set(configured
.split(',')
.map(port => parseInt(port.trim(), 10))
.filter(port => Number.isInteger(port) && port > 0 && port <= 65535)));
}
private getNetwork(): bitcoinjs.networks.Network {
const networkConfig = this.configService.get('NETWORK');
if (networkConfig === 'mainnet') {
return bitcoinjs.networks.bitcoin;
}
if (networkConfig === 'testnet') {
return bitcoinjs.networks.testnet;
}
if (networkConfig === 'regtest') {
return bitcoinjs.networks.regtest;
}
throw new Error('Invalid network configuration');
}
}
export class Sv2JobDeclarationConnection {
private readonly noiseSession: Sv2NoiseSession;
private readonly frameReader = new Sv2FrameReader(null);
private readonly frameWriter = new Sv2FrameWriter(null);
private handshakeBuffer = Buffer.alloc(0);
private handshakeComplete = false;
private destroyed = false;
private declareTxData = false;
private latestDeclaredJob: Sv2DeclaredMiningJob | null = null;
private readonly submittedSolutions = new Set<string>();
private readonly pendingDeclarations = new Map<number, {
job: Sv2DeclareMiningJob;
templateId: bigint;
template: ReturnType<TemplateProviderService['getLatestTemplate']>;
unknownTxPositionList: number[];
}>();
constructor(
private readonly socket: Socket,
stratumV2Service: StratumV2Service,
private readonly registry: Sv2JobDeclarationRegistryService,
private readonly customWorkService: CustomWorkService,
private readonly templateProvider: TemplateProviderService,
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly blocksService: BlocksService,
private readonly payoutSnapshotService: PayoutSnapshotService,
private readonly notificationService: NotificationService,
private readonly poolPayoutAddress: string,
private readonly network: bitcoinjs.networks.Network,
) {
this.noiseSession = new Sv2NoiseSession(stratumV2Service.getNoiseConfig());
}
public async start(): Promise<void> {
this.socket.setKeepAlive(true, 60_000);
this.socket.setNoDelay(true);
this.socket.on('data', data => void this.handleSocketData(data));
this.socket.on('error', () => this.close());
this.socket.on('close', () => this.destroyed = true);
}
private async handleSocketData(data: Buffer): Promise<void> {
if (this.destroyed) {
return;
}
try {
if (!this.handshakeComplete) {
await this.handleHandshakeData(data);
return;
}
await this.handleEncryptedData(data);
} catch (error) {
console.error(`[SV2 JDP] ${error.message}`);
this.close();
}
}
private async handleHandshakeData(data: Buffer): Promise<void> {
this.handshakeBuffer = Buffer.concat([this.handshakeBuffer, data]);
if (this.handshakeBuffer.length < SV2_NOISE_ACT1_SIZE) {
return;
}
const act1 = this.handshakeBuffer.subarray(0, SV2_NOISE_ACT1_SIZE);
const remainder = Buffer.from(this.handshakeBuffer.subarray(SV2_NOISE_ACT1_SIZE));
this.handshakeBuffer = Buffer.alloc(0);
await this.writeRaw(await this.noiseSession.processAct1(Buffer.from(act1)));
this.frameReader.setDecryptFn(ciphertext => this.noiseSession.decrypt(ciphertext));
this.frameWriter.setEncryptFn(plaintext => this.noiseSession.encrypt(plaintext));
this.handshakeComplete = true;
if (remainder.length > 0) {
await this.handleEncryptedData(remainder);
}
}
private async handleEncryptedData(data: Buffer): Promise<void> {
for (const frame of this.frameReader.feed(data)) {
switch (frame.header.msgType) {
case Sv2MsgType.SETUP_CONNECTION:
await this.handleSetupConnection(frame.payload);
break;
case Sv2MsgType.JDP_ALLOCATE_MINING_JOB_TOKEN:
await this.handleAllocateMiningJobToken(frame.payload);
break;
case Sv2MsgType.JDP_DECLARE_MINING_JOB:
await this.handleDeclareMiningJob(frame.payload);
break;
case Sv2MsgType.JDP_PROVIDE_MISSING_TRANSACTIONS_SUCCESS:
await this.handleProvideMissingTransactionsSuccess(frame.payload);
break;
case Sv2MsgType.JDP_PUSH_SOLUTION:
await this.handlePushSolution(frame.payload);
break;
default:
console.warn(`[SV2 JDP] Ignoring unsupported message type 0x${frame.header.msgType.toString(16)}`);
break;
}
}
}
private async handleSetupConnection(payload: Buffer): Promise<void> {
const setup = deserializeSetupConnection(new BufferReader(payload));
if (setup.protocol !== Sv2Protocol.JOB_DECLARATION) {
await this.sendFrame(Sv2MsgType.SETUP_CONNECTION_ERROR, serializeSetupConnectionError({
flags: 0,
errorCode: 'unsupported-protocol',
}));
this.close();
return;
}
this.declareTxData = (setup.flags & Sv2JdpSetupFlags.DECLARE_TX_DATA) !== 0;
await this.sendFrame(Sv2MsgType.SETUP_CONNECTION_SUCCESS, serializeSetupConnectionSuccess({
usedVersion: 2,
flags: 0,
}));
}
private async handleAllocateMiningJobToken(payload: Buffer): Promise<void> {
const request = deserializeAllocateMiningJobToken(new BufferReader(payload));
const token = this.registry.allocateToken(
request.userIdentifier,
this.buildPoolCoinbaseOutputs(request.userIdentifier),
);
await this.sendFrame(Sv2MsgType.JDP_ALLOCATE_MINING_JOB_TOKEN_SUCCESS, serializeAllocateMiningJobTokenSuccess({
requestId: request.requestId,
miningJobToken: token.token,
coinbaseOutputs: token.coinbaseOutputs,
}));
}
private async handleDeclareMiningJob(payload: Buffer): Promise<void> {
const job = deserializeDeclareMiningJob(new BufferReader(payload));
if (!this.declareTxData) {
await this.sendFrame(Sv2MsgType.JDP_DECLARE_MINING_JOB_ERROR, serializeDeclareMiningJobError({
requestId: job.requestId,
errorCode: 'declare-tx-data-not-negotiated',
errorDetails: Buffer.alloc(0),
}));
return;
}
try {
const validation = this.templateProvider.validateDeclaredWtxids({
version: job.version,
coinbaseTxPrefix: job.coinbaseTxPrefix,
wtxidList: job.wtxidList,
});
if (validation.errorCode === 'missing-transactions' && validation.template != null) {
this.pendingDeclarations.set(job.requestId, {
job,
templateId: validation.template.templateId,
template: validation.template,
unknownTxPositionList: validation.unknownTxPositionList,
});
await this.sendFrame(Sv2MsgType.JDP_PROVIDE_MISSING_TRANSACTIONS, serializeProvideMissingTransactions({
requestId: job.requestId,
unknownTxPositionList: validation.unknownTxPositionList,
}));
return;
}
if (!validation.valid) {
await this.sendFrame(Sv2MsgType.JDP_DECLARE_MINING_JOB_ERROR, serializeDeclareMiningJobError({
requestId: job.requestId,
errorCode: validation.errorCode ?? 'invalid-template',
errorDetails: Buffer.alloc(0),
}));
return;
}
const declared = this.registry.declareJob(job, {
templateId: validation.template.templateId,
template: validation.template,
validationMode: 'full_template',
});
this.latestDeclaredJob = declared;
await this.sendFrame(Sv2MsgType.JDP_DECLARE_MINING_JOB_SUCCESS, serializeDeclareMiningJobSuccess({
requestId: job.requestId,
newMiningJobToken: declared.token,
}));
} catch (error) {
await this.sendFrame(Sv2MsgType.JDP_DECLARE_MINING_JOB_ERROR, serializeDeclareMiningJobError({
requestId: job.requestId,
errorCode: error.message,
errorDetails: Buffer.alloc(0),
}));
}
}
private async handleProvideMissingTransactionsSuccess(payload: Buffer): Promise<void> {
const response = deserializeProvideMissingTransactionsSuccess(new BufferReader(payload));
const pending = this.pendingDeclarations.get(response.requestId);
if (pending == null) {
await this.sendFrame(Sv2MsgType.JDP_DECLARE_MINING_JOB_ERROR, serializeDeclareMiningJobError({
requestId: response.requestId,
errorCode: 'unknown-missing-transaction-request',
errorDetails: Buffer.alloc(0),
}));
return;
}
const validation = await this.templateProvider.validateProvidedTransactions({
expectedWtxids: pending.job.wtxidList,
unknownTxPositionList: pending.unknownTxPositionList,
transactionList: response.transactionList,
});
if (!validation.valid) {
this.pendingDeclarations.delete(response.requestId);
await this.sendFrame(Sv2MsgType.JDP_DECLARE_MINING_JOB_ERROR, serializeDeclareMiningJobError({
requestId: response.requestId,
errorCode: validation.errorCode ?? 'invalid-missing-transactions',
errorDetails: Buffer.alloc(0),
}));
return;
}
try {
const declared = this.registry.declareJob(pending.job, {
templateId: pending.templateId,
template: pending.template,
validationMode: 'full_template',
providedTransactions: response.transactionList,
});
this.latestDeclaredJob = declared;
this.pendingDeclarations.delete(response.requestId);
await this.sendFrame(Sv2MsgType.JDP_DECLARE_MINING_JOB_SUCCESS, serializeDeclareMiningJobSuccess({
requestId: pending.job.requestId,
newMiningJobToken: declared.token,
}));
} catch (error) {
this.pendingDeclarations.delete(response.requestId);
await this.sendFrame(Sv2MsgType.JDP_DECLARE_MINING_JOB_ERROR, serializeDeclareMiningJobError({
requestId: response.requestId,
errorCode: error.message,
errorDetails: Buffer.alloc(0),
}));
}
}
private async handlePushSolution(payload: Buffer): Promise<void> {
const solution = deserializePushSolution(new BufferReader(payload));
const declared = this.latestDeclaredJob;
if (declared == null) {
console.warn('[SV2 JDP] PushSolution received before a mining job was declared');
return;
}
const solutionKey = [
solution.prevHash.toString('hex'),
solution.extranonce.toString('hex'),
solution.nonce.toString(16),
solution.ntime.toString(16),
solution.nBits.toString(16),
solution.version.toString(16),
].join(':');
if (this.submittedSolutions.has(solutionKey)) {
return;
}
const template = declared.template
?? (declared.templateId == null
? this.templateProvider.getLatestTemplate()
: this.templateProvider.getTemplate(declared.templateId));
if (template == null) {
console.warn('[SV2 JDP] PushSolution skipped because the declared template is no longer available');
return;
}
try {
const block = this.templateProvider.buildBlockFromDeclaredJobSolution({
template,
job: declared.job,
providedTransactions: declared.providedTransactions,
solution,
});
const coinbaseValidation = this.templateProvider.validateCoinbaseTransactionHeight(
block.transactions[0].toBuffer(),
template.height,
);
if (!coinbaseValidation.valid) {
console.warn(`[SV2 JDP] PushSolution rejected locally: ${coinbaseValidation.errorCode}`);
return;
}
const blockHex = block.toHex(false);
this.submittedSolutions.add(solutionKey);
const result = await this.bitcoinRpcService.SUBMIT_BLOCK(blockHex);
const { address, worker } = this.parseUserIdentifier(declared.userIdentifier);
await this.blocksService.save({
height: template.height,
minerAddress: address,
worker,
sessionId: declared.token.toString('hex').slice(0, 8),
blockData: blockHex,
blockSubmissionResult: result,
payoutSnapshotId: template.jobTemplate.blockData.payoutSnapshotId ?? null,
});
await this.payoutSnapshotService.finalizeSnapshotForBlock({
payoutSnapshotId: template.jobTemplate.blockData.payoutSnapshotId,
blockHeight: template.height,
blockSubmissionResult: result,
});
await this.notificationService.notifySubscribersBlockFound(address, template.height, block, result);
console.log(`[SV2 JDP] PushSolution submitted block at height ${template.height}: ${result ?? 'accepted'}`);
} catch (error) {
console.error(`[SV2 JDP] PushSolution failed: ${error.message ?? error}`);
}
}
private buildPoolCoinbaseOutputs(userIdentifier: string): Buffer {
const address = this.poolPayoutAddress || userIdentifier.split('.')[0];
const script = bitcoinjs.address.toOutputScript(address, this.network);
const value = Buffer.alloc(8);
return Buffer.concat([
this.customWorkService.encodeBitcoinVarInt(1),
value,
this.customWorkService.encodeBitcoinVarInt(script.length),
script,
]);
}
private parseUserIdentifier(userIdentifier: string): { address: string; worker: string } {
const [address, ...workerParts] = userIdentifier.split('.');
return {
address: address || userIdentifier,
worker: workerParts.join('.') || 'jdp',
};
}
private async sendFrame(msgType: number, payload: Buffer): Promise<void> {
await this.writeRaw(this.frameWriter.writeFrame({
extensionType: 0,
msgType,
msgLength: payload.length,
}, payload));
}
private async writeRaw(data: Buffer): Promise<void> {
if (this.socket.destroyed || this.socket.writableEnded) {
return;
}
await new Promise<void>((resolve, reject) => {
this.socket.write(data, error => error ? reject(error) : resolve());
});
}
private close(): void {
this.destroyed = true;
if (!this.socket.destroyed) {
this.socket.destroy();
}
}
}
@@ -0,0 +1,290 @@
import * as bitcoinjs from 'bitcoinjs-lib';
import { BehaviorSubject, firstValueFrom } from 'rxjs';
import { MockRecording1 } from '../../test/models/MockRecording1';
import { BufferReader } from '../models/sv2/sv2-binary-codec';
import { MiningJob } from '../models/MiningJob';
import { Sv2MsgType, Sv2Protocol } from '../models/sv2/sv2-constants';
import {
deserializeSetupConnectionSuccess,
serializeSetupConnection,
} from '../models/sv2/sv2-messages';
import { serializeTdpCoinbaseOutputConstraints } from '../models/sv2/sv2-tdp-messages';
import {
deserializeTdpNewTemplate,
deserializeTdpRequestTransactionDataSuccess,
serializeTdpRequestTransactionData,
serializeTdpSubmitSolution,
} from '../models/sv2/sv2-tdp-messages';
import { StratumV1JobsService } from './stratum-v1-jobs.service';
import { Sv2TemplateDistributionConnection } from './sv2-template-distribution.service';
import { TemplateProviderService } from './template-provider.service';
describe('Sv2TemplateDistributionConnection compliance', () => {
it('waits for client CoinbaseOutputConstraints before sending templates', async () => {
const { connection, sentFrames, templateProvider, jobTemplate } = await createConnection();
await (connection as any).handleSetupConnection(serializeSetupConnection({
protocol: Sv2Protocol.TEMPLATE_DISTRIBUTION,
minVersion: 2,
maxVersion: 2,
flags: 0,
endpoint_host: 'localhost',
endpoint_port: 34265,
vendor: 'template-client',
hardwareVersion: '',
firmwareVersion: '',
deviceId: '',
}));
expect(sentFrames.map(frame => frame.msgType)).toEqual([Sv2MsgType.SETUP_CONNECTION_SUCCESS]);
const success = deserializeSetupConnectionSuccess(new BufferReader(sentFrames[0].payload));
expect(success.flags).toBe(0);
sentFrames.length = 0;
await (connection as any).handleFrame(
Sv2MsgType.TDP_COINBASE_OUTPUT_CONSTRAINTS,
serializeTdpCoinbaseOutputConstraints({
coinbaseOutputMaxAdditionalSize: 50_000,
coinbaseOutputMaxAdditionalSigops: 0,
}),
);
await Promise.resolve();
await Promise.resolve();
expect(sentFrames.some(frame => frame.msgType === Sv2MsgType.TDP_NEW_TEMPLATE)).toBe(true);
expect(sentFrames.some(frame => frame.msgType === Sv2MsgType.TDP_SET_NEW_PREV_HASH)).toBe(true);
const templateFrame = sentFrames.find(frame => frame.msgType === Sv2MsgType.TDP_NEW_TEMPLATE);
const template = deserializeTdpNewTemplate(new BufferReader(templateFrame.payload));
const expectedJob = new MiningJob(
bitcoinjs.networks.testnet,
jobTemplate.blockData.id,
[{ address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', percent: 100 }],
jobTemplate,
);
const expectedCoinbaseTx = expectedJob.cloneCoinbaseTransaction();
expect(template.coinbasePrefix).toEqual(expectedJob.getCoinbasePrefixBuffer());
expect(template.coinbasePrefix.includes(templateProvider.buildCoinbaseHeightPrefix(MockRecording1.BLOCK_TEMPLATE.height))).toBe(true);
expect(template.coinbaseTxVersion).toBe(expectedCoinbaseTx.version);
expect(template.coinbaseTxInputSequence).toBe(expectedCoinbaseTx.ins[0].sequence);
expect(template.coinbaseTxOutputsCount).toBe(expectedCoinbaseTx.outs.length);
expect(template.coinbaseTxOutputs).toEqual((connection as any).serializeCoinbaseOutputs(expectedCoinbaseTx));
expect(template.coinbaseTxLocktime).toBe(expectedCoinbaseTx.locktime);
});
it('serves transaction data for known templates', async () => {
const { connection, sentFrames, templateProvider, jobTemplate } = await createConnection();
const template = templateProvider.upsert(jobTemplate);
await (connection as any).handleFrame(
Sv2MsgType.TDP_REQUEST_TRANSACTION_DATA,
serializeTdpRequestTransactionData({ templateId: template.templateId }),
);
const successFrame = sentFrames.find(frame => frame.msgType === Sv2MsgType.TDP_REQUEST_TRANSACTION_DATA_SUCCESS);
const success = deserializeTdpRequestTransactionDataSuccess(new BufferReader(successFrame.payload));
expect(success.templateId).toBe(template.templateId);
expect(success.transactionList[0]).toEqual(Buffer.from(MockRecording1.BLOCK_TEMPLATE.transactions[0].data, 'hex'));
});
it('submits reconstructed blocks from SubmitSolution', async () => {
const { connection, templateProvider, jobTemplate, bitcoinRpcService, blocksService, payoutSnapshotService, notificationService } = await createConnection();
const template = templateProvider.upsert(jobTemplate);
const callOrder: string[] = [];
bitcoinRpcService.SUBMIT_BLOCK.mockImplementation(async () => {
callOrder.push('submit');
return null;
});
blocksService.save.mockImplementation(async () => {
callOrder.push('save');
});
const coinbaseTx = createCoinbaseTransaction(templateProvider.buildCoinbaseHeightPrefix(template.height));
await (connection as any).handleFrame(
Sv2MsgType.TDP_SUBMIT_SOLUTION,
serializeTdpSubmitSolution({
templateId: template.templateId,
version: template.version,
headerTimestamp: template.minNtime,
headerNonce: 123,
coinbaseTx: coinbaseTx.toBuffer(),
}),
);
expect(bitcoinRpcService.SUBMIT_BLOCK).toHaveBeenCalledWith(expect.stringMatching(/^[0-9a-f]+$/));
expect(blocksService.save).toHaveBeenCalledWith(expect.objectContaining({
height: MockRecording1.BLOCK_TEMPLATE.height,
minerAddress: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
worker: 'tdp',
blockData: expect.stringMatching(/^[0-9a-f]+$/),
}));
expect(callOrder).toEqual(['submit', 'save']);
expect(payoutSnapshotService.finalizeSnapshotForBlock).toHaveBeenCalledWith({
payoutSnapshotId: jobTemplate.blockData.payoutSnapshotId,
blockHeight: MockRecording1.BLOCK_TEMPLATE.height,
blockSubmissionResult: null,
});
expect(notificationService.notifySubscribersBlockFound).toHaveBeenCalledWith(
'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
MockRecording1.BLOCK_TEMPLATE.height,
expect.any(bitcoinjs.Block),
null,
);
});
it('does not resubmit duplicate SubmitSolution messages', async () => {
const { connection, templateProvider, jobTemplate, bitcoinRpcService } = await createConnection();
const template = templateProvider.upsert(jobTemplate);
const coinbaseTx = createCoinbaseTransaction(templateProvider.buildCoinbaseHeightPrefix(template.height));
const payload = serializeTdpSubmitSolution({
templateId: template.templateId,
version: template.version,
headerTimestamp: template.minNtime,
headerNonce: 123,
coinbaseTx: coinbaseTx.toBuffer(),
});
await (connection as any).handleFrame(Sv2MsgType.TDP_SUBMIT_SOLUTION, payload);
await (connection as any).handleFrame(Sv2MsgType.TDP_SUBMIT_SOLUTION, payload);
expect(bitcoinRpcService.SUBMIT_BLOCK).toHaveBeenCalledTimes(1);
});
it('rejects SubmitSolution locally when the coinbase is missing the BIP34 height', async () => {
const warn = jest.spyOn(console, 'warn').mockImplementation(() => undefined);
const { connection, templateProvider, jobTemplate, bitcoinRpcService, blocksService } = await createConnection();
const template = templateProvider.upsert(jobTemplate);
const coinbaseTx = createCoinbaseTransaction(Buffer.from('/public-pool-sri-jdc-e2e/', 'utf8'));
await (connection as any).handleFrame(
Sv2MsgType.TDP_SUBMIT_SOLUTION,
serializeTdpSubmitSolution({
templateId: template.templateId,
version: template.version,
headerTimestamp: template.minNtime,
headerNonce: 123,
coinbaseTx: coinbaseTx.toBuffer(),
}),
);
expect(bitcoinRpcService.SUBMIT_BLOCK).not.toHaveBeenCalled();
expect(blocksService.save).not.toHaveBeenCalled();
expect(warn).toHaveBeenCalledWith(expect.stringContaining('invalid-job-param-value-coinbase_tx_prefix'));
warn.mockRestore();
});
it('submits using a served template snapshot after provider cache cleanup', async () => {
const { connection, sentFrames, templateProvider, jobTemplate, bitcoinRpcService, blocksService } = await createConnection();
await (connection as any).sendTemplate(jobTemplate);
const templateId = BigInt(parseInt(jobTemplate.blockData.id, 16));
jest.spyOn(templateProvider, 'getTemplate').mockReturnValue(undefined);
sentFrames.length = 0;
const coinbaseTx = createCoinbaseTransaction(templateProvider.buildCoinbaseHeightPrefix(jobTemplate.blockData.height));
await (connection as any).handleFrame(
Sv2MsgType.TDP_SUBMIT_SOLUTION,
serializeTdpSubmitSolution({
templateId,
version: jobTemplate.block.version,
headerTimestamp: jobTemplate.block.timestamp,
headerNonce: 123,
coinbaseTx: coinbaseTx.toBuffer(),
}),
);
expect(bitcoinRpcService.SUBMIT_BLOCK).toHaveBeenCalledWith(expect.stringMatching(/^[0-9a-f]+$/));
expect(blocksService.save).toHaveBeenCalledWith(expect.objectContaining({
height: MockRecording1.BLOCK_TEMPLATE.height,
blockSubmissionResult: null,
}));
expect(sentFrames.some(frame => frame.msgType === Sv2MsgType.TDP_REQUEST_TRANSACTION_DATA_ERROR)).toBe(false);
});
async function createConnection(): Promise<{
connection: Sv2TemplateDistributionConnection;
sentFrames: any[];
templateProvider: TemplateProviderService;
jobTemplate: any;
bitcoinRpcService: { SUBMIT_BLOCK: jest.Mock };
blocksService: { save: jest.Mock };
payoutSnapshotService: { finalizeSnapshotForBlock: jest.Mock };
notificationService: { notifySubscribersBlockFound: jest.Mock };
}> {
const blockTemplate$ = new BehaviorSubject(MockRecording1.BLOCK_TEMPLATE);
const jobsService = new StratumV1JobsService({
newBlockTemplate$: blockTemplate$.asObservable(),
miningInfo: { blocks: MockRecording1.BLOCK_TEMPLATE.height },
} as any);
const jobTemplate = await firstValueFrom(jobsService.newMiningJob$);
const templateProvider = new TemplateProviderService(jobsService);
const bitcoinRpcService = {
SUBMIT_BLOCK: jest.fn().mockResolvedValue(null),
};
const blocksService = {
save: jest.fn().mockResolvedValue(undefined),
};
const payoutSnapshotService = {
finalizeSnapshotForBlock: jest.fn().mockResolvedValue({ finalized: false, reason: 'disabled' }),
};
const notificationService = {
notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined),
};
const socket = {
setKeepAlive: jest.fn(),
setNoDelay: jest.fn(),
on: jest.fn(),
write: jest.fn((data, callback) => callback?.()),
destroy: jest.fn(),
destroyed: false,
writableEnded: false,
};
const connection = new Sv2TemplateDistributionConnection(
socket as any,
{
getNoiseConfig: () => ({
staticKeypair: {
privateKey: Buffer.alloc(32),
publicKey: Buffer.alloc(64),
},
certificateMessage: {
version: 0,
validFrom: 0,
notValidAfter: 0,
signature: Buffer.alloc(64),
},
}),
} as any,
jobsService,
templateProvider,
bitcoinRpcService as any,
blocksService as any,
payoutSnapshotService as any,
notificationService as any,
'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
bitcoinjs.networks.testnet,
);
const sentFrames: any[] = [];
(connection as any).sendFrame = jest.fn((msgType: number, payload: Buffer) => {
sentFrames.push({ msgType, payload });
return Promise.resolve();
});
return {
connection,
sentFrames,
templateProvider,
jobTemplate,
bitcoinRpcService,
blocksService,
payoutSnapshotService,
notificationService,
};
}
function createCoinbaseTransaction(script: Buffer): bitcoinjs.Transaction {
const tx = new bitcoinjs.Transaction();
tx.version = 2;
tx.addInput(Buffer.alloc(32), 0xffffffff, 0xffffffff, script);
tx.addOutput(Buffer.from('6a', 'hex'), 0);
return tx;
}
});
@@ -0,0 +1,444 @@
import { Injectable, OnModuleInit } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import * as bitcoinjs from 'bitcoinjs-lib';
import { Server, Socket } from 'net';
import { Subscription } from 'rxjs';
import { BlocksService } from '../ORM/blocks/blocks.service';
import { PayoutSnapshotService } from '../ORM/payout-snapshot/payout-snapshot.service';
import { SV2_NOISE_ACT1_SIZE, Sv2MsgType, Sv2Protocol } from '../models/sv2/sv2-constants';
import { Sv2FrameReader, Sv2FrameWriter } from '../models/sv2/sv2-frame';
import { BufferReader } from '../models/sv2/sv2-binary-codec';
import {
deserializeSetupConnection,
serializeSetupConnectionError,
serializeSetupConnectionSuccess,
} from '../models/sv2/sv2-messages';
import { AddressObject, MiningJob } from '../models/MiningJob';
import { Sv2NoiseSession } from '../models/sv2/sv2-noise';
import {
deserializeTdpCoinbaseOutputConstraints,
deserializeTdpRequestTransactionData,
deserializeTdpSubmitSolution,
serializeTdpNewTemplate,
serializeTdpRequestTransactionDataError,
serializeTdpRequestTransactionDataSuccess,
serializeTdpSetNewPrevHash,
} from '../models/sv2/sv2-tdp-messages';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { NotificationService } from './notification.service';
import { IJobTemplate, StratumV1JobsService } from './stratum-v1-jobs.service';
import { StratumV2Service } from './stratum-v2.service';
import { TemplateProviderService, TemplateProviderTemplate } from './template-provider.service';
interface TdpCoinbaseTemplateFields {
coinbaseTxVersion: number;
coinbasePrefix: Buffer;
coinbaseTxInputSequence: number;
coinbaseTxValueRemaining: bigint;
coinbaseTxOutputsCount: number;
coinbaseTxOutputs: Buffer;
coinbaseTxLocktime: number;
}
@Injectable()
export class Sv2TemplateDistributionService implements OnModuleInit {
private readonly servers: Server[] = [];
constructor(
private readonly configService: ConfigService,
private readonly stratumV2Service: StratumV2Service,
private readonly jobsService: StratumV1JobsService,
private readonly templateProvider: TemplateProviderService,
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly blocksService: BlocksService,
private readonly payoutSnapshotService: PayoutSnapshotService,
private readonly notificationService: NotificationService,
) {}
public async onModuleInit(): Promise<void> {
if (process.env.API_ONLY === 'true' || process.env.MASTER === 'true') {
return;
}
const ports = this.getPorts();
if (ports.length === 0) {
return;
}
await this.stratumV2Service.ensureInitialized();
for (const port of ports) {
this.startServer(port);
}
}
private startServer(port: number): void {
const server = new Server(socket => {
void new Sv2TemplateDistributionConnection(
socket,
this.stratumV2Service,
this.jobsService,
this.templateProvider,
this.bitcoinRpcService,
this.blocksService,
this.payoutSnapshotService,
this.notificationService,
this.getPoolPayoutAddress(),
this.getNetwork(),
).start();
});
server.on('error', error => console.error(`SV2 TDP server error on port ${port}: ${error.message}`));
server.listen(port, () => console.log(`SV2 Template Distribution server is listening on port ${port}`));
this.servers.push(server);
}
private getPorts(): number[] {
const configured = this.configService.get<string>('SV2_TDP_PORTS');
if (!configured?.trim()) {
return [];
}
return Array.from(new Set(configured
.split(',')
.map(port => parseInt(port.trim(), 10))
.filter(port => Number.isInteger(port) && port > 0 && port <= 65535)));
}
private getPoolPayoutAddress(): string {
return this.configService.get<string>('SV2_TDP_POOL_PAYOUT_ADDRESS')
|| this.configService.get<string>('SV2_JDP_POOL_PAYOUT_ADDRESS')
|| this.configService.get<string>('DATUM_POOL_PAYOUT_ADDRESS')
|| this.configService.get<string>('PAYOUT_FEE_ADDRESS')
|| this.configService.get<string>('DEV_FEE_ADDRESS')
|| '';
}
private getNetwork(): bitcoinjs.networks.Network {
const networkConfig = this.configService.get('NETWORK');
if (networkConfig === 'mainnet') {
return bitcoinjs.networks.bitcoin;
}
if (networkConfig === 'testnet') {
return bitcoinjs.networks.testnet;
}
if (networkConfig === 'regtest') {
return bitcoinjs.networks.regtest;
}
throw new Error('Invalid network configuration');
}
}
export class Sv2TemplateDistributionConnection {
private readonly noiseSession: Sv2NoiseSession;
private readonly frameReader = new Sv2FrameReader(null);
private readonly frameWriter = new Sv2FrameWriter(null);
private handshakeBuffer = Buffer.alloc(0);
private handshakeComplete = false;
private destroyed = false;
private subscription: Subscription | null = null;
private readonly submittedSolutions = new Set<string>();
private readonly servedTemplates = new Map<string, TemplateProviderTemplate>();
constructor(
private readonly socket: Socket,
stratumV2Service: StratumV2Service,
private readonly jobsService: StratumV1JobsService,
private readonly templateProvider: TemplateProviderService,
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly blocksService: BlocksService,
private readonly payoutSnapshotService: PayoutSnapshotService,
private readonly notificationService: NotificationService,
private readonly poolPayoutAddress: string,
private readonly network: bitcoinjs.networks.Network,
) {
this.noiseSession = new Sv2NoiseSession(stratumV2Service.getNoiseConfig());
}
public async start(): Promise<void> {
this.socket.setKeepAlive(true, 60_000);
this.socket.setNoDelay(true);
this.socket.on('data', data => void this.handleSocketData(data));
this.socket.on('error', () => this.close());
this.socket.on('close', () => this.close());
}
private async handleSocketData(data: Buffer): Promise<void> {
if (this.destroyed) {
return;
}
try {
if (!this.handshakeComplete) {
await this.handleHandshakeData(data);
return;
}
for (const frame of this.frameReader.feed(data)) {
await this.handleFrame(frame.header.msgType, frame.payload);
}
} catch (error) {
console.error(`[SV2 TDP] ${error.message}`);
this.close();
}
}
private async handleHandshakeData(data: Buffer): Promise<void> {
this.handshakeBuffer = Buffer.concat([this.handshakeBuffer, data]);
if (this.handshakeBuffer.length < SV2_NOISE_ACT1_SIZE) {
return;
}
const act1 = this.handshakeBuffer.subarray(0, SV2_NOISE_ACT1_SIZE);
const remainder = Buffer.from(this.handshakeBuffer.subarray(SV2_NOISE_ACT1_SIZE));
this.handshakeBuffer = Buffer.alloc(0);
await this.writeRaw(await this.noiseSession.processAct1(Buffer.from(act1)));
this.frameReader.setDecryptFn(ciphertext => this.noiseSession.decrypt(ciphertext));
this.frameWriter.setEncryptFn(plaintext => this.noiseSession.encrypt(plaintext));
this.handshakeComplete = true;
if (remainder.length > 0) {
for (const frame of this.frameReader.feed(remainder)) {
await this.handleFrame(frame.header.msgType, frame.payload);
}
}
}
private async handleFrame(msgType: number, payload: Buffer): Promise<void> {
switch (msgType) {
case Sv2MsgType.SETUP_CONNECTION:
await this.handleSetupConnection(payload);
break;
case Sv2MsgType.TDP_COINBASE_OUTPUT_CONSTRAINTS:
deserializeTdpCoinbaseOutputConstraints(new BufferReader(payload));
this.subscribeTemplates();
break;
case Sv2MsgType.TDP_REQUEST_TRANSACTION_DATA: {
const request = deserializeTdpRequestTransactionData(new BufferReader(payload));
await this.handleRequestTransactionData(request.templateId);
break;
}
case Sv2MsgType.TDP_SUBMIT_SOLUTION: {
const solution = deserializeTdpSubmitSolution(new BufferReader(payload));
await this.handleSubmitSolution(solution);
break;
}
default:
break;
}
}
private async handleSetupConnection(payload: Buffer): Promise<void> {
const setup = deserializeSetupConnection(new BufferReader(payload));
if (setup.protocol !== Sv2Protocol.TEMPLATE_DISTRIBUTION) {
await this.sendFrame(Sv2MsgType.SETUP_CONNECTION_ERROR, serializeSetupConnectionError({
flags: 0,
errorCode: 'unsupported-protocol',
}));
this.close();
return;
}
await this.sendFrame(Sv2MsgType.SETUP_CONNECTION_SUCCESS, serializeSetupConnectionSuccess({
usedVersion: 2,
flags: 0,
}));
}
private subscribeTemplates(): void {
if (this.subscription != null) {
return;
}
this.subscription = this.jobsService.newMiningJob$.subscribe({
next: template => void this.sendTemplate(template),
error: () => this.close(),
});
}
private async sendTemplate(template: IJobTemplate): Promise<void> {
const cachedTemplate = this.templateProvider.upsert(template);
const templateId = BigInt(parseInt(template.blockData.id, 16));
this.servedTemplates.set(templateId.toString(), cachedTemplate);
const merklePath = template.merkle_branch.map(branch => Buffer.from(branch, 'hex'));
const coinbase = this.buildCoinbaseTemplateFields(template);
await this.sendFrame(Sv2MsgType.TDP_NEW_TEMPLATE, serializeTdpNewTemplate({
templateId,
futureTemplate: true,
version: template.block.version,
...coinbase,
merklePath,
}));
await this.sendFrame(Sv2MsgType.TDP_SET_NEW_PREV_HASH, serializeTdpSetNewPrevHash({
templateId,
prevHash: Buffer.from(template.block.prevHash),
headerTimestamp: template.block.timestamp,
nBits: template.block.bits,
target: cachedTemplate.target,
}));
}
private buildCoinbaseTemplateFields(template: IJobTemplate): TdpCoinbaseTemplateFields {
const job = new MiningJob(
this.network,
template.blockData.id,
this.getPayoutInformation(template),
template,
);
const coinbaseTx = job.cloneCoinbaseTransaction();
return {
coinbaseTxVersion: coinbaseTx.version,
coinbasePrefix: job.getCoinbasePrefixBuffer(),
coinbaseTxInputSequence: coinbaseTx.ins[0]?.sequence ?? 0xffffffff,
coinbaseTxValueRemaining: BigInt(template.blockData.coinbasevalue),
coinbaseTxOutputsCount: coinbaseTx.outs.length,
coinbaseTxOutputs: this.serializeCoinbaseOutputs(coinbaseTx),
coinbaseTxLocktime: coinbaseTx.locktime,
};
}
private getPayoutInformation(template: IJobTemplate): AddressObject[] {
if (template.blockData.payoutOutputs?.length > 0) {
return template.blockData.payoutOutputs;
}
return [{ address: this.poolPayoutAddress, percent: 100 }];
}
private serializeCoinbaseOutputs(coinbaseTx: bitcoinjs.Transaction): Buffer {
return Buffer.concat(coinbaseTx.outs.map(output => {
const value = Buffer.alloc(8);
value.writeBigUInt64LE(BigInt(output.value), 0);
return Buffer.concat([
value,
this.encodeBitcoinVarInt(output.script.length),
Buffer.from(output.script),
]);
}));
}
private encodeBitcoinVarInt(value: number): Buffer {
if (!Number.isSafeInteger(value) || value < 0) {
throw new RangeError(`Invalid Bitcoin varint value ${value}`);
}
if (value < 0xfd) {
return Buffer.from([value]);
}
if (value <= 0xffff) {
const result = Buffer.alloc(3);
result[0] = 0xfd;
result.writeUInt16LE(value, 1);
return result;
}
if (value <= 0xffffffff) {
const result = Buffer.alloc(5);
result[0] = 0xfe;
result.writeUInt32LE(value, 1);
return result;
}
const result = Buffer.alloc(9);
result[0] = 0xff;
result.writeBigUInt64LE(BigInt(value), 1);
return result;
}
private async handleRequestTransactionData(templateId: bigint): Promise<void> {
const template = this.getServedOrProviderTemplate(templateId);
if (template == null) {
await this.sendFrame(Sv2MsgType.TDP_REQUEST_TRANSACTION_DATA_ERROR, serializeTdpRequestTransactionDataError({
templateId,
errorCode: 'template-not-found',
}));
return;
}
this.servedTemplates.set(templateId.toString(), template);
await this.sendFrame(Sv2MsgType.TDP_REQUEST_TRANSACTION_DATA_SUCCESS, serializeTdpRequestTransactionDataSuccess({
templateId,
excessData: Buffer.alloc(0),
transactionList: template.transactionList,
}));
}
private async handleSubmitSolution(solution: ReturnType<typeof deserializeTdpSubmitSolution>): Promise<void> {
const template = this.getServedOrProviderTemplate(solution.templateId);
if (template == null) {
await this.sendFrame(Sv2MsgType.TDP_REQUEST_TRANSACTION_DATA_ERROR, serializeTdpRequestTransactionDataError({
templateId: solution.templateId,
errorCode: 'template-not-found',
}));
return;
}
const coinbaseValidation = this.templateProvider.validateCoinbaseTransactionHeight(solution.coinbaseTx, template.height);
if (!coinbaseValidation.valid) {
console.warn(`[SV2 TDP] SubmitSolution rejected locally: ${coinbaseValidation.errorCode}`);
return;
}
const block = this.templateProvider.buildBlockFromSolution({
template,
coinbaseTx: solution.coinbaseTx,
version: solution.version,
headerTimestamp: solution.headerTimestamp,
headerNonce: solution.headerNonce,
});
const blockHex = block.toHex(false);
const solutionKey = `${solution.templateId.toString()}:${block.getId()}`;
if (this.submittedSolutions.has(solutionKey)) {
return;
}
this.submittedSolutions.add(solutionKey);
const result = await this.bitcoinRpcService.SUBMIT_BLOCK(blockHex);
await this.blocksService.save({
height: template.height,
minerAddress: this.poolPayoutAddress || 'sv2-tdp',
worker: 'tdp',
sessionId: solution.templateId.toString(16).slice(-8).padStart(8, '0'),
blockData: blockHex,
blockSubmissionResult: result,
payoutSnapshotId: template.jobTemplate.blockData.payoutSnapshotId ?? null,
});
await this.payoutSnapshotService.finalizeSnapshotForBlock({
payoutSnapshotId: template.jobTemplate.blockData.payoutSnapshotId,
blockHeight: template.height,
blockSubmissionResult: result,
});
await this.notificationService.notifySubscribersBlockFound(
this.poolPayoutAddress || 'sv2-tdp',
template.height,
block,
result,
);
if (result != null && result !== 'SUCCESS!') {
console.warn(`[SV2 TDP] SubmitSolution rejected: ${result}`);
}
}
private getServedOrProviderTemplate(templateId: bigint): TemplateProviderTemplate | undefined {
return this.servedTemplates.get(templateId.toString())
?? this.templateProvider.getTemplate(templateId);
}
private async sendFrame(msgType: number, payload: Buffer): Promise<void> {
await this.writeRaw(this.frameWriter.writeFrame({
extensionType: 0,
msgType,
msgLength: payload.length,
}, payload));
}
private async writeRaw(data: Buffer): Promise<void> {
if (this.socket.destroyed || this.socket.writableEnded) {
return;
}
await new Promise<void>((resolve, reject) => {
this.socket.write(data, error => error ? reject(error) : resolve());
});
}
private close(): void {
if (this.destroyed) {
return;
}
this.destroyed = true;
this.subscription?.unsubscribe();
this.subscription = null;
if (!this.socket.destroyed) {
this.socket.destroy();
}
}
}
@@ -0,0 +1,331 @@
import { BehaviorSubject, firstValueFrom } from 'rxjs';
import * as bitcoinjs from 'bitcoinjs-lib';
import { MockRecording1 } from '../../test/models/MockRecording1';
import { CustomWorkService } from './custom-work.service';
import { StratumV1JobsService } from './stratum-v1-jobs.service';
import { TemplateProviderService } from './template-provider.service';
describe('TemplateProviderService', () => {
it('indexes template transactions and validates known wtxids', async () => {
const { provider, jobTemplate } = await createProvider();
const template = provider.upsert(jobTemplate);
const firstTx = MockRecording1.BLOCK_TEMPLATE.transactions[0];
expect(template.transactionList[0]).toEqual(Buffer.from(firstTx.data, 'hex'));
expect(provider.getTemplate(template.templateId)).toBe(template);
const validation = provider.validateDeclaredWtxids({
version: MockRecording1.BLOCK_TEMPLATE.version,
wtxidList: [Buffer.from(firstTx.hash, 'hex')],
});
expect(validation.valid).toBe(true);
expect(validation.template).toBe(template);
});
it('reports unknown transaction positions and validates supplied missing tx bytes', async () => {
const { provider, jobTemplate } = await createProvider();
provider.upsert(jobTemplate);
const firstTx = MockRecording1.BLOCK_TEMPLATE.transactions[0];
const unknownWtxid = Buffer.alloc(32, 0xaa);
const validation = provider.validateDeclaredWtxids({
version: MockRecording1.BLOCK_TEMPLATE.version,
wtxidList: [Buffer.from(firstTx.hash, 'hex'), unknownWtxid],
});
expect(validation.valid).toBe(false);
expect(validation.errorCode).toBe('missing-transactions');
expect(validation.unknownTxPositionList).toEqual([1]);
const missingTxValidation = await provider.validateProvidedTransactions({
expectedWtxids: [Buffer.from(firstTx.hash, 'hex'), unknownWtxid],
unknownTxPositionList: [1],
transactionList: [Buffer.from(firstTx.data, 'hex')],
});
expect(missingTxValidation.valid).toBe(false);
expect(missingTxValidation.errorCode).toBe('transaction-wtxid-mismatch');
});
it('rejects duplicate declared wtxids before requesting missing transactions', async () => {
const { provider, jobTemplate } = await createProvider();
provider.upsert(jobTemplate);
const firstTx = MockRecording1.BLOCK_TEMPLATE.transactions[0];
const validation = provider.validateDeclaredWtxids({
version: MockRecording1.BLOCK_TEMPLATE.version,
wtxidList: [
Buffer.from(firstTx.hash, 'hex'),
Buffer.from(firstTx.hash, 'hex'),
],
});
expect(validation.valid).toBe(false);
expect(validation.errorCode).toBe('duplicate-transactions');
});
it('builds and validates the BIP34 coinbase height prefix for SV2 templates', async () => {
const { provider, jobTemplate } = await createProvider();
const template = provider.upsert(jobTemplate);
const heightPrefix = provider.buildCoinbaseHeightPrefix(template.height);
const coinbaseTx = createCoinbaseTransaction(heightPrefix);
expect(heightPrefix).toEqual(Buffer.concat([
Buffer.from([bitcoinjs.script.number.encode(template.height).length]),
bitcoinjs.script.number.encode(template.height),
]));
expect(provider.validateCoinbaseTransactionHeight(coinbaseTx.toBuffer(), template.height)).toEqual({ valid: true });
});
it('rejects declared SV2 JDP coinbase prefixes missing the required BIP34 height', async () => {
const { provider, jobTemplate } = await createProvider();
const template = provider.upsert(jobTemplate);
const badScript = Buffer.from('/public-pool-sri-jdc-e2e/', 'utf8');
const badPrefix = createCoinbasePrefixForScript(badScript, badScript);
const validation = provider.validateDeclaredWtxids({
version: MockRecording1.BLOCK_TEMPLATE.version,
coinbaseTxPrefix: badPrefix,
wtxidList: [],
});
expect(validation).toMatchObject({
valid: false,
errorCode: 'invalid-job-param-value-coinbase_tx_prefix',
template,
});
});
it('accepts declared SV2 JDP coinbase prefixes that start with the required BIP34 height', async () => {
const { provider, jobTemplate } = await createProvider();
const template = provider.upsert(jobTemplate);
const heightPrefix = provider.buildCoinbaseHeightPrefix(template.height);
const validPrefix = createCoinbasePrefixForScript(
Buffer.concat([heightPrefix, Buffer.from('/public-pool-sri-jdc-e2e/', 'utf8')]),
heightPrefix,
);
const validation = provider.validateDeclaredWtxids({
version: MockRecording1.BLOCK_TEMPLATE.version,
coinbaseTxPrefix: validPrefix,
wtxidList: [],
});
expect(validation.valid).toBe(true);
expect(validation.template).toBe(template);
});
it('accepts DATUM template metadata matching the latest template fast path', async () => {
const { provider, jobTemplate } = await createProvider();
const template = provider.upsert(jobTemplate);
const nBits = Buffer.alloc(4);
nBits.writeUInt32LE(template.nBits, 0);
const validation = provider.validateDatumTemplateFastPath({
prevBlockHash: Buffer.from(template.prevHash),
nBits,
height: template.height,
version: template.version ^ 0x2000,
coinbaseValue: template.coinbaseValue,
totalWeight: template.weightLimit,
totalSize: template.sizeLimit,
totalSigops: template.sigopLimit,
merkleBranches: [Buffer.alloc(32)],
});
expect(validation.valid).toBe(true);
expect(validation.template).toBe(template);
});
it('rejects stale or impossible DATUM template metadata in the fast path', async () => {
const { provider, jobTemplate } = await createProvider();
const template = provider.upsert(jobTemplate);
const nBits = Buffer.alloc(4);
nBits.writeUInt32LE(template.nBits, 0);
expect(provider.validateDatumTemplateFastPath({
prevBlockHash: Buffer.alloc(32, 0xaa),
nBits,
height: template.height,
})).toMatchObject({ valid: false, errorCode: 'prevhash-mismatch' });
expect(provider.validateDatumTemplateFastPath({
prevBlockHash: Buffer.from(template.prevHash),
nBits,
height: template.height + 1,
})).toMatchObject({ valid: false, errorCode: 'height-mismatch' });
expect(provider.validateDatumTemplateFastPath({
prevBlockHash: Buffer.from(template.prevHash),
nBits,
height: template.height,
coinbaseValue: template.coinbaseValue + 1n,
})).toMatchObject({ valid: false, errorCode: 'coinbase-value-too-high' });
expect(provider.validateDatumTemplateFastPath({
prevBlockHash: Buffer.from(template.prevHash),
nBits,
height: template.height,
totalWeight: template.weightLimit + 1,
})).toMatchObject({ valid: false, errorCode: 'weight-limit-exceeded' });
});
it('validates supplied missing transactions against Bitcoin Core mempool policy', async () => {
const firstTx = MockRecording1.BLOCK_TEMPLATE.transactions[0];
const { provider, bitcoinRpcService } = await createProvider({
TEST_MEMPOOL_ACCEPT: jest.fn().mockResolvedValue([{ allowed: true }]),
});
const validation = await provider.validateProvidedTransactions({
expectedWtxids: [Buffer.from(firstTx.hash, 'hex')],
unknownTxPositionList: [0],
transactionList: [Buffer.from(firstTx.data, 'hex')],
});
expect(validation.valid).toBe(true);
expect(bitcoinRpcService.TEST_MEMPOOL_ACCEPT).toHaveBeenCalledWith([Buffer.from(firstTx.data, 'hex')]);
});
it('rejects malformed supplied missing transactions before mempool policy checks', async () => {
const firstTx = MockRecording1.BLOCK_TEMPLATE.transactions[0];
const { provider, bitcoinRpcService } = await createProvider({
TEST_MEMPOOL_ACCEPT: jest.fn().mockResolvedValue([{ allowed: true }]),
});
const validation = await provider.validateProvidedTransactions({
expectedWtxids: [Buffer.from(firstTx.hash, 'hex')],
unknownTxPositionList: [0],
transactionList: [Buffer.from('010203', 'hex')],
});
expect(validation.valid).toBe(false);
expect(validation.errorCode).toMatch(/^transaction-decode-failed:0:/);
expect(bitcoinRpcService.TEST_MEMPOOL_ACCEPT).not.toHaveBeenCalled();
});
it('validates generic transaction data for SV2 and DATUM callers', async () => {
const firstTx = MockRecording1.BLOCK_TEMPLATE.transactions[0];
const { provider } = await createProvider();
const valid = provider.validateTransactionData({
transactionList: [Buffer.from(firstTx.data, 'hex')],
expectedCount: 1,
maxTotalBytes: Buffer.from(firstTx.data, 'hex').length,
});
expect(valid.valid).toBe(true);
expect(valid.transactions?.[0].wtxid).toBe(firstTx.hash);
expect(provider.validateTransactionData({
transactionList: [Buffer.from(firstTx.data, 'hex')],
expectedCount: 2,
})).toMatchObject({
valid: false,
errorCode: 'transaction-count-mismatch:2:1',
});
expect(provider.validateTransactionData({
transactionList: [Buffer.from(firstTx.data, 'hex')],
maxTotalBytes: 1,
})).toMatchObject({
valid: false,
errorCode: expect.stringMatching(/^transaction-bytes-exceed-limit:/),
});
});
it('validates DATUM transaction blobs against the advertised coinbase merkle path', async () => {
const { provider } = await createProvider();
const coinbaseTx = createTestTransaction(0);
const transactionList = [
createTestTransaction(1),
createTestTransaction(2),
createTestTransaction(3),
];
const merklePath = provider.buildCoinbaseMerklePath(transactionList);
const merkleRoot = new CustomWorkService().computeMerkleRoot(coinbaseTx, merklePath);
const block = new bitcoinjs.Block();
block.transactions = [
bitcoinjs.Transaction.fromBuffer(coinbaseTx),
...transactionList.map(tx => bitcoinjs.Transaction.fromBuffer(tx)),
];
expect(merkleRoot).toEqual(bitcoinjs.Block.calculateMerkleRoot(block.transactions, false));
expect(provider.validateTransactionData({
transactionList,
expectedCount: transactionList.length,
expectedCoinbaseMerklePath: merklePath,
})).toMatchObject({ valid: true });
const wrongPath = merklePath.map(branch => Buffer.from(branch));
wrongPath[0][0] ^= 0xff;
expect(provider.validateTransactionData({
transactionList,
expectedCount: transactionList.length,
expectedCoinbaseMerklePath: wrongPath,
})).toMatchObject({
valid: false,
errorCode: 'merkle-path-mismatch:0',
});
});
it('rejects supplied missing transactions rejected by Bitcoin Core mempool policy', async () => {
const firstTx = MockRecording1.BLOCK_TEMPLATE.transactions[0];
const { provider } = await createProvider({
TEST_MEMPOOL_ACCEPT: jest.fn().mockResolvedValue([{
allowed: false,
rejectReason: 'txn-mempool-conflict',
}]),
});
const validation = await provider.validateProvidedTransactions({
expectedWtxids: [Buffer.from(firstTx.hash, 'hex')],
unknownTxPositionList: [0],
transactionList: [Buffer.from(firstTx.data, 'hex')],
});
expect(validation.valid).toBe(false);
expect(validation.errorCode).toBe('txn-mempool-conflict');
});
async function createProvider(bitcoinRpcService?: any): Promise<{
provider: TemplateProviderService;
jobTemplate: any;
bitcoinRpcService: any;
}> {
const blockTemplate$ = new BehaviorSubject(MockRecording1.BLOCK_TEMPLATE);
const jobsService = new StratumV1JobsService({
newBlockTemplate$: blockTemplate$.asObservable(),
miningInfo: { blocks: MockRecording1.BLOCK_TEMPLATE.height },
} as any);
const jobTemplate = await firstValueFrom(jobsService.newMiningJob$);
return {
provider: new TemplateProviderService(jobsService, bitcoinRpcService),
jobTemplate,
bitcoinRpcService,
};
}
function createTestTransaction(seed: number): Buffer {
const tx = new bitcoinjs.Transaction();
tx.version = 2;
tx.addInput(Buffer.alloc(32, seed), 0xffffffff, 0xffffffff, Buffer.from([seed & 0xff]));
tx.addOutput(Buffer.from('6a', 'hex'), 0);
return tx.toBuffer();
}
function createCoinbaseTransaction(script: Buffer): bitcoinjs.Transaction {
const tx = new bitcoinjs.Transaction();
tx.version = 2;
tx.addInput(Buffer.alloc(32), 0xffffffff, 0xffffffff, script);
tx.addOutput(Buffer.from('6a', 'hex'), 0);
return tx;
}
function createCoinbasePrefixForScript(fullScript: Buffer, prefixScript: Buffer): Buffer {
const tx = createCoinbaseTransaction(fullScript);
const serialized = tx.toBuffer();
const scriptStart = serialized.indexOf(fullScript);
if (scriptStart < 0) {
throw new Error('test coinbase script not found');
}
return serialized.subarray(0, scriptStart + prefixScript.length);
}
});
+656
View File
@@ -0,0 +1,656 @@
import { Injectable, OnModuleInit } from '@nestjs/common';
import * as bitcoinjs from 'bitcoinjs-lib';
import { Subscription } from 'rxjs';
import { IBlockTemplateTx } from '../models/bitcoin-rpc/IBlockTemplate';
import { Sv2DeclareMiningJob, Sv2PushSolution } from '../models/sv2/sv2-jdp-messages';
import { hash256 } from '../utils/hash.utils';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { IJobTemplate, StratumV1JobsService } from './stratum-v1-jobs.service';
export interface TemplateProviderTransaction {
position: number;
data: Buffer;
txid: string;
wtxid: string;
fee: number;
sigops: number;
weight: number;
}
export interface TemplateProviderTemplate {
templateId: bigint;
jobTemplate: IJobTemplate;
version: number;
prevHash: Buffer;
nBits: number;
minNtime: number;
coinbaseValue: bigint;
networkDifficulty: number;
target: Buffer;
height: number;
sigopLimit: number;
sizeLimit: number;
weightLimit: number;
transactions: TemplateProviderTransaction[];
transactionList: Buffer[];
txidMap: Map<string, TemplateProviderTransaction>;
wtxidMap: Map<string, TemplateProviderTransaction>;
createdAt: number;
}
export interface DeclareMiningJobValidationResult {
valid: boolean;
errorCode?: string;
template?: TemplateProviderTemplate;
unknownTxPositionList?: number[];
}
export interface DatumTemplateValidationInput {
prevBlockHash?: Buffer;
nBits?: Buffer;
height?: number;
version?: number;
coinbaseValue?: bigint;
totalWeight?: number;
totalSize?: number;
totalSigops?: number;
merkleBranches?: Buffer[];
}
export interface TemplateFastPathValidationResult {
valid: boolean;
errorCode?: string;
template?: TemplateProviderTemplate;
}
export interface TransactionDataValidationResult {
valid: boolean;
errorCode?: string;
transactions?: TemplateProviderTransaction[];
totalBytes?: number;
}
export interface CoinbaseHeightValidationResult {
valid: boolean;
errorCode?: string;
}
const VERSION_ROLLING_MASK = 0x1fffe000;
@Injectable()
export class TemplateProviderService implements OnModuleInit {
private readonly templates = new Map<string, TemplateProviderTemplate>();
private readonly retentionMs = this.readPositiveInt('SV2_TEMPLATE_RETENTION_MS', 10 * 60 * 1000);
private subscription: Subscription | null = null;
private latestTemplateId: string | null = null;
constructor(
private readonly jobsService: StratumV1JobsService,
private readonly bitcoinRpcService?: BitcoinRpcService,
) {}
public onModuleInit(): void {
if (this.subscription != null) {
return;
}
this.subscription = this.jobsService.newMiningJob$.subscribe(template => {
this.upsert(template);
});
}
public upsert(jobTemplate: IJobTemplate): TemplateProviderTemplate {
this.cleanup();
const templateId = BigInt(parseInt(jobTemplate.blockData.id, 16));
const transactions = this.buildTransactions(jobTemplate.blockData.transactions ?? []);
const template: TemplateProviderTemplate = {
templateId,
jobTemplate,
version: jobTemplate.block.version,
prevHash: Buffer.from(jobTemplate.block.prevHash),
nBits: jobTemplate.block.bits,
minNtime: jobTemplate.block.timestamp,
coinbaseValue: BigInt(jobTemplate.blockData.coinbasevalue),
networkDifficulty: jobTemplate.blockData.networkDifficulty,
target: this.targetFromDifficulty(jobTemplate.blockData.networkDifficulty),
height: jobTemplate.blockData.height,
sigopLimit: jobTemplate.blockData.sigoplimit ?? 80_000,
sizeLimit: jobTemplate.blockData.sizelimit ?? 4_000_000,
weightLimit: jobTemplate.blockData.weightlimit ?? 4_000_000,
transactions,
transactionList: transactions.map(tx => Buffer.from(tx.data)),
txidMap: this.buildMap(transactions, tx => tx.txid),
wtxidMap: this.buildMap(transactions, tx => tx.wtxid),
createdAt: Date.now(),
};
this.templates.set(template.templateId.toString(), template);
this.latestTemplateId = template.templateId.toString();
return template;
}
public getTemplate(templateId: bigint | number | string): TemplateProviderTemplate | undefined {
this.cleanup();
return this.templates.get(BigInt(templateId).toString());
}
public getLatestTemplate(): TemplateProviderTemplate | undefined {
this.cleanup();
return this.latestTemplateId == null ? undefined : this.templates.get(this.latestTemplateId);
}
public validateDeclaredWtxids(input: {
version: number;
wtxidList: Buffer[];
coinbaseTxPrefix?: Buffer;
}): DeclareMiningJobValidationResult {
const template = this.getLatestTemplate();
if (template == null) {
return { valid: false, errorCode: 'template-not-found' };
}
if (input.version !== template.version) {
return { valid: false, errorCode: 'version-mismatch', template };
}
if (input.coinbaseTxPrefix != null) {
const coinbaseHeight = this.validateDeclaredCoinbasePrefixHeight(input.coinbaseTxPrefix, template.height);
if (!coinbaseHeight.valid) {
return { valid: false, errorCode: coinbaseHeight.errorCode, template };
}
}
const seen = new Set<string>();
const unknownTxPositionList: number[] = [];
let duplicateFound = false;
input.wtxidList.forEach((wtxid, index) => {
const aliases = this.hashAliases(wtxid.toString('hex'));
if (aliases.some(alias => seen.has(alias))) {
duplicateFound = true;
unknownTxPositionList.push(index);
return;
}
aliases.forEach(alias => seen.add(alias));
if (!this.lookupWtxid(template, wtxid)) {
unknownTxPositionList.push(index);
}
});
if (duplicateFound) {
return {
valid: false,
errorCode: 'duplicate-transactions',
template,
unknownTxPositionList,
};
}
return {
valid: unknownTxPositionList.length === 0,
errorCode: unknownTxPositionList.length > 0 ? 'missing-transactions' : undefined,
template,
unknownTxPositionList,
};
}
public buildCoinbaseHeightPrefix(height: number): Buffer {
const encodedHeight = bitcoinjs.script.number.encode(height);
return Buffer.concat([Buffer.from([encodedHeight.length]), encodedHeight]);
}
public validateCoinbaseTransactionHeight(coinbaseTx: Buffer, height: number): CoinbaseHeightValidationResult {
let tx: bitcoinjs.Transaction;
try {
tx = bitcoinjs.Transaction.fromBuffer(coinbaseTx);
} catch (error) {
return {
valid: false,
errorCode: `invalid-coinbase-transaction:${error instanceof Error ? error.message : String(error)}`,
};
}
if (tx.ins.length !== 1) {
return { valid: false, errorCode: 'invalid-coinbase-input-count' };
}
return this.validateCoinbaseScriptHeight(tx.ins[0].script, height);
}
public validateDeclaredCoinbasePrefixHeight(coinbaseTxPrefix: Buffer, height: number): CoinbaseHeightValidationResult {
const scriptStart = this.getCoinbaseInputScriptStart(coinbaseTxPrefix);
if (scriptStart == null) {
return { valid: false, errorCode: 'invalid-job-param-value-coinbase_tx_prefix' };
}
const expected = this.buildCoinbaseHeightPrefix(height);
if (coinbaseTxPrefix.length < scriptStart + expected.length) {
return { valid: false, errorCode: 'invalid-job-param-value-coinbase_tx_prefix' };
}
return this.validateCoinbaseScriptHeight(coinbaseTxPrefix.subarray(scriptStart), height);
}
public async validateProvidedTransactions(input: {
expectedWtxids: Buffer[];
unknownTxPositionList: number[];
transactionList: Buffer[];
}): Promise<{ valid: boolean; errorCode?: string }> {
if (input.unknownTxPositionList.length !== input.transactionList.length) {
return { valid: false, errorCode: 'missing-transaction-count-mismatch' };
}
const transactionValidation = this.validateTransactionData({
transactionList: input.transactionList,
expectedCount: input.transactionList.length,
});
if (!transactionValidation.valid) {
return {
valid: false,
errorCode: transactionValidation.errorCode,
};
}
for (let i = 0; i < input.transactionList.length; i++) {
const position = input.unknownTxPositionList[i];
const expected = input.expectedWtxids[position];
const provided = transactionValidation.transactions![i];
if (expected == null || !this.matchesHash(expected, provided.wtxid)) {
return { valid: false, errorCode: 'transaction-wtxid-mismatch' };
}
}
const mempoolResult = await this.validateMempoolPolicy(input.transactionList);
if (!mempoolResult.valid) {
return mempoolResult;
}
return { valid: true };
}
public validateTransactionData(input: {
transactionList: Buffer[];
expectedCount?: number;
maxTotalBytes?: number;
expectedCoinbaseMerklePath?: Buffer[];
}): TransactionDataValidationResult {
if (input.expectedCount != null && input.transactionList.length !== input.expectedCount) {
return {
valid: false,
errorCode: `transaction-count-mismatch:${input.expectedCount}:${input.transactionList.length}`,
};
}
let totalBytes = 0;
const transactions: TemplateProviderTransaction[] = [];
for (let i = 0; i < input.transactionList.length; i++) {
const rawTransaction = input.transactionList[i];
totalBytes += rawTransaction.length;
if (input.maxTotalBytes != null && totalBytes > input.maxTotalBytes) {
return {
valid: false,
errorCode: `transaction-bytes-exceed-limit:${totalBytes}:${input.maxTotalBytes}`,
totalBytes,
};
}
try {
transactions.push(this.parseTransaction(rawTransaction, i));
} catch (error) {
return {
valid: false,
errorCode: `transaction-decode-failed:${i}:${error instanceof Error ? error.message : String(error)}`,
totalBytes,
};
}
}
if (input.expectedCoinbaseMerklePath != null) {
const merklePath = this.buildCoinbaseMerklePath(input.transactionList);
if (merklePath.length !== input.expectedCoinbaseMerklePath.length) {
return {
valid: false,
errorCode: `merkle-path-length-mismatch:${input.expectedCoinbaseMerklePath.length}:${merklePath.length}`,
totalBytes,
};
}
for (let i = 0; i < merklePath.length; i++) {
if (!merklePath[i].equals(input.expectedCoinbaseMerklePath[i])) {
return {
valid: false,
errorCode: `merkle-path-mismatch:${i}`,
totalBytes,
};
}
}
}
return {
valid: true,
transactions,
totalBytes,
};
}
public buildCoinbaseMerklePath(transactionList: Buffer[]): Buffer[] {
const coinbasePlaceholderHash = Buffer.alloc(32);
let level = [
coinbasePlaceholderHash,
...transactionList.map(transaction => bitcoinjs.Transaction.fromBuffer(transaction).getHash(false)),
];
let index = 0;
const path: Buffer[] = [];
while (level.length > 1) {
const siblingIndex = index ^ 1;
path.push(Buffer.from(level[siblingIndex] ?? level[index]));
const nextLevel: Buffer[] = [];
for (let i = 0; i < level.length; i += 2) {
const left = level[i];
const right = level[i + 1] ?? left;
nextLevel.push(hash256(Buffer.concat([left, right])));
}
level = nextLevel;
index = Math.floor(index / 2);
}
return path;
}
public validateDatumTemplateFastPath(input: DatumTemplateValidationInput): TemplateFastPathValidationResult {
const template = this.getLatestTemplate();
if (template == null) {
return { valid: false, errorCode: 'template-not-found' };
}
if (input.prevBlockHash != null && !this.matchesHash(input.prevBlockHash, template.prevHash.toString('hex'))) {
return { valid: false, errorCode: 'prevhash-mismatch', template };
}
if (input.nBits != null) {
if (input.nBits.length !== 4) {
return { valid: false, errorCode: 'nbits-mismatch', template };
}
if (input.nBits.readUInt32LE(0) !== template.nBits) {
return { valid: false, errorCode: 'nbits-mismatch', template };
}
}
if (input.height != null && input.height !== template.height) {
return { valid: false, errorCode: 'height-mismatch', template };
}
if (input.version != null && !this.isCompatibleVersion(input.version, template.version)) {
return { valid: false, errorCode: 'version-mismatch', template };
}
if (input.coinbaseValue != null && input.coinbaseValue > template.coinbaseValue) {
return { valid: false, errorCode: 'coinbase-value-too-high', template };
}
if (input.totalWeight != null && input.totalWeight > template.weightLimit) {
return { valid: false, errorCode: 'weight-limit-exceeded', template };
}
if (input.totalSize != null && input.totalSize > template.sizeLimit) {
return { valid: false, errorCode: 'size-limit-exceeded', template };
}
if (input.totalSigops != null && input.totalSigops > template.sigopLimit) {
return { valid: false, errorCode: 'sigop-limit-exceeded', template };
}
if (input.merkleBranches != null && input.merkleBranches.some(branch => branch.length !== 32)) {
return { valid: false, errorCode: 'bad-merkle-branch', template };
}
return { valid: true, template };
}
public buildBlockFromSolution(input: {
template: TemplateProviderTemplate;
coinbaseTx: Buffer;
version: number;
headerTimestamp: number;
headerNonce: number;
}): bitcoinjs.Block {
const block = new bitcoinjs.Block();
block.version = input.version;
block.prevHash = Buffer.from(input.template.prevHash);
block.timestamp = input.headerTimestamp;
block.bits = input.template.nBits;
block.nonce = input.headerNonce;
block.transactions = [
bitcoinjs.Transaction.fromBuffer(input.coinbaseTx),
...input.template.transactionList.map(tx => bitcoinjs.Transaction.fromBuffer(tx)),
];
block.merkleRoot = bitcoinjs.Block.calculateMerkleRoot(block.transactions, false);
try {
block.witnessCommit = bitcoinjs.Block.calculateMerkleRoot(block.transactions, true);
} catch {
block.witnessCommit = null;
}
return block;
}
public buildBlockFromDeclaredJobSolution(input: {
template: TemplateProviderTemplate;
job: Sv2DeclareMiningJob;
providedTransactions?: Buffer[];
solution: Sv2PushSolution;
}): bitcoinjs.Block {
const transactionMap = this.buildDeclaredTransactionMap(input.template, input.providedTransactions ?? []);
const coinbaseTx = Buffer.concat([
input.job.coinbaseTxPrefix,
input.solution.extranonce,
input.job.coinbaseTxSuffix,
]);
const transactions = [bitcoinjs.Transaction.fromBuffer(coinbaseTx)];
for (const wtxid of input.job.wtxidList) {
const transaction = this.lookupDeclaredTransaction(transactionMap, wtxid);
if (transaction == null) {
throw new Error('declared-transaction-not-found');
}
transactions.push(bitcoinjs.Transaction.fromBuffer(transaction));
}
const block = new bitcoinjs.Block();
block.version = input.solution.version;
block.prevHash = Buffer.from(input.template.prevHash);
block.timestamp = input.solution.ntime;
block.bits = input.solution.nBits;
block.nonce = input.solution.nonce;
block.transactions = transactions;
block.merkleRoot = bitcoinjs.Block.calculateMerkleRoot(block.transactions, false);
try {
block.witnessCommit = bitcoinjs.Block.calculateMerkleRoot(block.transactions, true);
} catch {
block.witnessCommit = null;
}
return block;
}
private buildTransactions(rawTransactions: IBlockTemplateTx[]): TemplateProviderTransaction[] {
return rawTransactions.map((tx, index) => {
const parsed = this.parseTransaction(Buffer.from(tx.data, 'hex'), index);
return {
...parsed,
txid: tx.txid || parsed.txid,
wtxid: tx.hash || parsed.wtxid,
fee: tx.fee ?? 0,
sigops: tx.sigops ?? 0,
weight: tx.weight ?? 0,
};
});
}
private parseTransaction(data: Buffer, position: number): TemplateProviderTransaction {
const tx = bitcoinjs.Transaction.fromBuffer(data);
return {
position,
data: Buffer.from(data),
txid: tx.getId(),
wtxid: Buffer.from(tx.getHash(true)).reverse().toString('hex'),
fee: 0,
sigops: 0,
weight: 0,
};
}
private validateCoinbaseScriptHeight(script: Buffer, height: number): CoinbaseHeightValidationResult {
const expected = this.buildCoinbaseHeightPrefix(height);
if (script.length < expected.length || !script.subarray(0, expected.length).equals(expected)) {
return { valid: false, errorCode: 'invalid-job-param-value-coinbase_tx_prefix' };
}
return { valid: true };
}
private getCoinbaseInputScriptStart(txPrefix: Buffer): number | null {
let cursor = 4;
if (txPrefix.length < cursor + 1) {
return null;
}
if (txPrefix.length >= 6 && txPrefix[4] === 0x00 && txPrefix[5] !== 0x00) {
cursor = 6;
}
const inputCount = this.readBitcoinVarInt(txPrefix, cursor);
if (inputCount == null || inputCount.value !== 1n) {
return null;
}
cursor = inputCount.offset;
if (txPrefix.length < cursor + 36) {
return null;
}
cursor += 36;
const scriptLength = this.readBitcoinVarInt(txPrefix, cursor);
if (scriptLength == null) {
return null;
}
return scriptLength.offset;
}
private readBitcoinVarInt(buffer: Buffer, offset: number): { value: bigint; offset: number } | null {
if (offset >= buffer.length) {
return null;
}
const first = buffer[offset];
if (first < 0xfd) {
return { value: BigInt(first), offset: offset + 1 };
}
if (first === 0xfd) {
if (offset + 3 > buffer.length) {
return null;
}
return { value: BigInt(buffer.readUInt16LE(offset + 1)), offset: offset + 3 };
}
if (first === 0xfe) {
if (offset + 5 > buffer.length) {
return null;
}
return { value: BigInt(buffer.readUInt32LE(offset + 1)), offset: offset + 5 };
}
if (offset + 9 > buffer.length) {
return null;
}
return { value: buffer.readBigUInt64LE(offset + 1), offset: offset + 9 };
}
private buildDeclaredTransactionMap(
template: TemplateProviderTemplate,
providedTransactions: Buffer[],
): Map<string, Buffer> {
const map = new Map<string, Buffer>();
for (const tx of template.transactions) {
for (const alias of this.hashAliases(tx.wtxid)) {
map.set(alias, Buffer.from(tx.data));
}
}
for (const provided of providedTransactions) {
const parsed = this.parseTransaction(provided, -1);
for (const alias of this.hashAliases(parsed.wtxid)) {
map.set(alias, Buffer.from(provided));
}
}
return map;
}
private lookupDeclaredTransaction(map: Map<string, Buffer>, hash: Buffer): Buffer | undefined {
for (const alias of this.hashAliases(hash.toString('hex'))) {
const found = map.get(alias);
if (found != null) {
return found;
}
}
return undefined;
}
private buildMap(
transactions: TemplateProviderTransaction[],
selector: (tx: TemplateProviderTransaction) => string,
): Map<string, TemplateProviderTransaction> {
const map = new Map<string, TemplateProviderTransaction>();
for (const tx of transactions) {
for (const key of this.hashAliases(selector(tx))) {
map.set(key, tx);
}
}
return map;
}
private lookupWtxid(template: TemplateProviderTemplate, hash: Buffer): TemplateProviderTransaction | undefined {
for (const alias of this.hashAliases(hash.toString('hex'))) {
const found = template.wtxidMap.get(alias);
if (found != null) {
return found;
}
}
return undefined;
}
private matchesHash(candidate: Buffer, expectedHex: string): boolean {
const aliases = this.hashAliases(expectedHex);
return aliases.includes(candidate.toString('hex').toLowerCase());
}
private isCompatibleVersion(candidate: number, templateVersion: number): boolean {
return (((candidate >>> 0) ^ (templateVersion >>> 0)) & ~VERSION_ROLLING_MASK) === 0;
}
private async validateMempoolPolicy(transactionList: Buffer[]): Promise<{ valid: boolean; errorCode?: string }> {
if (transactionList.length === 0 || this.bitcoinRpcService?.TEST_MEMPOOL_ACCEPT == null) {
return { valid: true };
}
try {
const results = await this.bitcoinRpcService.TEST_MEMPOOL_ACCEPT(transactionList);
const rejected = results.find(result => !result.allowed);
if (rejected != null) {
return {
valid: false,
errorCode: rejected.rejectReason || rejected.rejectDetails || 'mempool-rejected-transaction',
};
}
return { valid: true };
} catch (error) {
return {
valid: false,
errorCode: `mempool-validation-failed:${error.message ?? error}`,
};
}
}
private hashAliases(hashHex: string): string[] {
const normalized = hashHex.toLowerCase();
const reversed = Buffer.from(normalized, 'hex').reverse().toString('hex');
return normalized === reversed ? [normalized] : [normalized, reversed];
}
private targetFromDifficulty(difficulty: number): Buffer {
const maxTarget = (BigInt(0xffff) << BigInt(8 * (0x1d - 3)));
const divisor = BigInt(Math.max(1, Math.floor(difficulty)));
const target = maxTarget / divisor;
const result = Buffer.alloc(32);
let remaining = target;
for (let i = 0; i < 32; i++) {
result[i] = Number(remaining & 0xffn);
remaining >>= 8n;
}
return result;
}
private cleanup(): void {
const cutoff = Date.now() - this.retentionMs;
for (const [id, template] of this.templates.entries()) {
if (template.createdAt < cutoff) {
this.templates.delete(id);
if (this.latestTemplateId === id) {
this.latestTemplateId = null;
}
}
}
}
private readPositiveInt(name: string, fallback: number): number {
const parsed = parseInt(process.env[name] ?? '', 10);
return Number.isFinite(parsed) && parsed > 0 ? parsed : fallback;
}
}