mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 17:15:03 -07:00
Stratum V2 initial commit
This commit is contained in:
@@ -11,6 +11,7 @@ describe('StratumV1Service', () => {
|
||||
|
||||
let service: StratumV1Service;
|
||||
let clientService;
|
||||
let stratumV2Service;
|
||||
let consoleLogSpy: jest.SpyInstance;
|
||||
let consoleWarnSpy: jest.SpyInstance;
|
||||
|
||||
@@ -19,6 +20,10 @@ describe('StratumV1Service', () => {
|
||||
clientService = {
|
||||
deleteAll: jest.fn().mockResolvedValue(undefined)
|
||||
};
|
||||
stratumV2Service = {
|
||||
ensureInitialized: jest.fn().mockResolvedValue(undefined),
|
||||
createClient: jest.fn()
|
||||
};
|
||||
service = new StratumV1Service(
|
||||
{} as any,
|
||||
clientService,
|
||||
@@ -27,7 +32,8 @@ describe('StratumV1Service', () => {
|
||||
{} as any,
|
||||
{} as any,
|
||||
{} as any,
|
||||
{} as any
|
||||
{} as any,
|
||||
stratumV2Service as any
|
||||
);
|
||||
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
|
||||
consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined);
|
||||
@@ -145,6 +151,29 @@ describe('StratumV1Service', () => {
|
||||
expect((service as any).getTlsHandshakeTimeoutMs()).toBe(5000);
|
||||
});
|
||||
|
||||
it('should detect JSON-RPC as Stratum V1', () => {
|
||||
const firstChunk = Buffer.from('{"id":1,"method":"mining.subscribe","params":[]}\n');
|
||||
|
||||
expect((service as any).detectProtocol(firstChunk)).toBe('v1');
|
||||
});
|
||||
|
||||
it('should detect binary Noise traffic as Stratum V2', () => {
|
||||
const firstChunk = Buffer.concat([
|
||||
Buffer.from([0x01, 0x02, 0x03, 0x04]),
|
||||
Buffer.alloc(60, 0xaa)
|
||||
]);
|
||||
|
||||
expect((service as any).detectProtocol(firstChunk)).toBe('v2');
|
||||
});
|
||||
|
||||
it('should reject recognizable TLS client hello on the unified plain stratum port', () => {
|
||||
expect((service as any).detectProtocol(Buffer.from([0x16, 0x03, 0x01]))).toBeNull();
|
||||
});
|
||||
|
||||
it('should route binary data that only shares a TLS first byte to Stratum V2', () => {
|
||||
expect((service as any).detectProtocol(Buffer.from([0x16, 0xaa, 0xbb]))).toBe('v2');
|
||||
});
|
||||
|
||||
function restoreEnv(key: string, value: string | undefined) {
|
||||
if (value == null) {
|
||||
delete process.env[key];
|
||||
|
||||
@@ -4,6 +4,7 @@ import { Server, Socket } from 'net';
|
||||
import { monitorEventLoopDelay } from 'perf_hooks';
|
||||
|
||||
import { StratumV1Client } from '../models/StratumV1Client';
|
||||
import { StratumV2Client } from '../models/StratumV2Client';
|
||||
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
|
||||
import { BlocksService } from '../ORM/blocks/blocks.service';
|
||||
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
|
||||
@@ -11,6 +12,7 @@ import { ClientService } from '../ORM/client/client.service';
|
||||
import { BitcoinRpcService } from './bitcoin-rpc.service';
|
||||
import { NotificationService } from './notification.service';
|
||||
import { StratumV1JobsService } from './stratum-v1-jobs.service';
|
||||
import { StratumV2Service } from './stratum-v2.service';
|
||||
|
||||
import { readFileSync } from 'fs';
|
||||
import { TlsOptions, TLSSocket, createServer } from 'tls';
|
||||
@@ -54,7 +56,8 @@ export class StratumV1Service implements OnModuleInit {
|
||||
private readonly blocksService: BlocksService,
|
||||
private readonly configService: ConfigService,
|
||||
private readonly stratumV1JobsService: StratumV1JobsService,
|
||||
private readonly addressSettingsService: AddressSettingsService
|
||||
private readonly addressSettingsService: AddressSettingsService,
|
||||
private readonly stratumV2Service: StratumV2Service
|
||||
) {
|
||||
|
||||
}
|
||||
@@ -107,21 +110,12 @@ export class StratumV1Service implements OnModuleInit {
|
||||
// Set 15-minute timeout
|
||||
socket.setTimeout(1000 * 60 * 15);
|
||||
|
||||
const client = new StratumV1Client(
|
||||
socket,
|
||||
this.stratumV1JobsService,
|
||||
this.bitcoinRpcService,
|
||||
this.clientService,
|
||||
this.clientStatisticsService,
|
||||
this.notificationService,
|
||||
this.blocksService,
|
||||
this.configService,
|
||||
this.addressSettingsService
|
||||
);
|
||||
let client: StratumV1Client | StratumV2Client = null;
|
||||
let protocol: 'v1' | 'v2' | null = null;
|
||||
|
||||
// Unified cleanup function
|
||||
const cleanup = async (reason: string) => {
|
||||
if (client.extraNonceAndSessionId != null) {
|
||||
if (client != null && (protocol === 'v2' || (client as StratumV1Client).extraNonceAndSessionId != null)) {
|
||||
await client.destroy();
|
||||
if (reason == 'Error') {
|
||||
this.errorClosure++;
|
||||
@@ -155,8 +149,30 @@ export class StratumV1Service implements OnModuleInit {
|
||||
await cleanup("Error");
|
||||
});
|
||||
|
||||
//
|
||||
socket.once('data', async (firstChunk: Buffer) => {
|
||||
try {
|
||||
protocol = this.detectProtocol(firstChunk);
|
||||
if (protocol === 'v1') {
|
||||
client = this.createV1Client(socket);
|
||||
socket.emit('data', firstChunk);
|
||||
return;
|
||||
}
|
||||
|
||||
if (protocol === 'v2') {
|
||||
await this.stratumV2Service.ensureInitialized();
|
||||
client = this.stratumV2Service.createClient(socket, firstChunk);
|
||||
return;
|
||||
}
|
||||
|
||||
if (!socket.destroyed) {
|
||||
socket.end();
|
||||
socket.destroy();
|
||||
}
|
||||
} catch (error) {
|
||||
console.error(`Protocol detection failed: ${error.message}`);
|
||||
await cleanup('Error');
|
||||
}
|
||||
});
|
||||
|
||||
});
|
||||
|
||||
@@ -169,6 +185,20 @@ export class StratumV1Service implements OnModuleInit {
|
||||
return server;
|
||||
}
|
||||
|
||||
private createV1Client(socket: Socket): StratumV1Client {
|
||||
return new StratumV1Client(
|
||||
socket,
|
||||
this.stratumV1JobsService,
|
||||
this.bitcoinRpcService,
|
||||
this.clientService,
|
||||
this.clientStatisticsService,
|
||||
this.notificationService,
|
||||
this.blocksService,
|
||||
this.configService,
|
||||
this.addressSettingsService
|
||||
);
|
||||
}
|
||||
|
||||
private startSecureSocketServer(port: number) {
|
||||
const listener: StratumListenerState = {
|
||||
port,
|
||||
@@ -196,17 +226,7 @@ export class StratumV1Service implements OnModuleInit {
|
||||
// Set 15-minute timeout
|
||||
socket.setTimeout(1000 * 60 * 15);
|
||||
|
||||
const client = new StratumV1Client(
|
||||
socket,
|
||||
this.stratumV1JobsService,
|
||||
this.bitcoinRpcService,
|
||||
this.clientService,
|
||||
this.clientStatisticsService,
|
||||
this.notificationService,
|
||||
this.blocksService,
|
||||
this.configService,
|
||||
this.addressSettingsService
|
||||
);
|
||||
const client = this.createV1Client(socket);
|
||||
|
||||
const cleanup = async (reason: string) => {
|
||||
if (client.extraNonceAndSessionId != null) {
|
||||
@@ -389,6 +409,70 @@ export class StratumV1Service implements OnModuleInit {
|
||||
return this.getPositiveIntegerEnv('STRATUM_TLS_HANDSHAKE_TIMEOUT_MS', DEFAULT_TLS_HANDSHAKE_TIMEOUT_MS);
|
||||
}
|
||||
|
||||
private detectProtocol(firstChunk: Buffer): 'v1' | 'v2' | null {
|
||||
if (firstChunk.length === 0) {
|
||||
return null;
|
||||
}
|
||||
|
||||
if (this.looksLikeJsonRpc(firstChunk)) {
|
||||
return 'v1';
|
||||
}
|
||||
|
||||
// TLS ClientHello. Secure SV1 remains on SECURE_STRATUM_PORTS.
|
||||
if (this.looksLikeTlsClientHello(firstChunk)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
// HTTP on a stratum port is not supported in this branch.
|
||||
if (this.looksLikeHttpRequest(firstChunk)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return 'v2';
|
||||
}
|
||||
|
||||
private looksLikeJsonRpc(firstChunk: Buffer): boolean {
|
||||
let firstNonWhitespace = -1;
|
||||
for (let i = 0; i < firstChunk.length; i++) {
|
||||
const byte = firstChunk[i];
|
||||
if (byte === 0x20 || byte === 0x09 || byte === 0x0a || byte === 0x0d) {
|
||||
continue;
|
||||
}
|
||||
firstNonWhitespace = i;
|
||||
break;
|
||||
}
|
||||
|
||||
if (firstNonWhitespace < 0 || firstChunk[firstNonWhitespace] !== 0x7b) {
|
||||
return false;
|
||||
}
|
||||
|
||||
for (const byte of firstChunk) {
|
||||
const isWhitespace = byte === 0x09 || byte === 0x0a || byte === 0x0d;
|
||||
const isPrintableAscii = byte >= 0x20 && byte <= 0x7e;
|
||||
if (!isWhitespace && !isPrintableAscii) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
private looksLikeTlsClientHello(firstChunk: Buffer): boolean {
|
||||
return firstChunk.length >= 3
|
||||
&& firstChunk[0] === 0x16
|
||||
&& firstChunk[1] === 0x03;
|
||||
}
|
||||
|
||||
private looksLikeHttpRequest(firstChunk: Buffer): boolean {
|
||||
const prefix = firstChunk.subarray(0, Math.min(firstChunk.length, 8)).toString('ascii').toUpperCase();
|
||||
return prefix.startsWith('GET ')
|
||||
|| prefix.startsWith('POST ')
|
||||
|| prefix.startsWith('PUT ')
|
||||
|| prefix.startsWith('PATCH ')
|
||||
|| prefix.startsWith('HEAD ')
|
||||
|| prefix.startsWith('OPTIONS ');
|
||||
}
|
||||
|
||||
private getPositiveIntegerEnv(key: string, fallback: number) {
|
||||
const configured = parseInt(process.env[key], 10);
|
||||
if (Number.isFinite(configured) && configured > 0) {
|
||||
|
||||
@@ -0,0 +1,211 @@
|
||||
import { Injectable, OnModuleInit } from '@nestjs/common';
|
||||
import { ConfigService } from '@nestjs/config';
|
||||
import * as crypto from 'crypto';
|
||||
import { Server, Socket } from 'net';
|
||||
|
||||
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
|
||||
import { BlocksService } from '../ORM/blocks/blocks.service';
|
||||
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
|
||||
import { ClientService } from '../ORM/client/client.service';
|
||||
import { StratumV2Client } from '../models/StratumV2Client';
|
||||
import { encodeSv2AuthorityPublicKey } from '../models/sv2/sv2-authority-key';
|
||||
import { Sv2ExtranonceManager } from '../models/sv2/sv2-extranonce-manager';
|
||||
import {
|
||||
createSignatureNoiseMessage,
|
||||
generateServerKeypair,
|
||||
Sv2NoiseConfig,
|
||||
Sv2ServerKeypair,
|
||||
xOnlyPubKeyFromPriv,
|
||||
} from '../models/sv2/sv2-noise';
|
||||
import { BitcoinRpcService } from './bitcoin-rpc.service';
|
||||
import { NotificationService } from './notification.service';
|
||||
import { StratumV1JobsService } from './stratum-v1-jobs.service';
|
||||
|
||||
@Injectable()
|
||||
export class StratumV2Service implements OnModuleInit {
|
||||
private readonly servers: Server[] = [];
|
||||
private authorityPrivKey: Buffer;
|
||||
private authorityPublicKeyXOnly: Buffer;
|
||||
private authorityKeyConfigured = false;
|
||||
private serverKeypair: Sv2ServerKeypair;
|
||||
private noiseConfig: Sv2NoiseConfig;
|
||||
private channelIdCounter = 1;
|
||||
private readonly extranonceManager = new Sv2ExtranonceManager();
|
||||
|
||||
constructor(
|
||||
private readonly bitcoinRpcService: BitcoinRpcService,
|
||||
private readonly clientService: ClientService,
|
||||
private readonly clientStatisticsService: ClientStatisticsService,
|
||||
private readonly notificationService: NotificationService,
|
||||
private readonly blocksService: BlocksService,
|
||||
private readonly configService: ConfigService,
|
||||
private readonly stratumV1JobsService: StratumV1JobsService,
|
||||
private readonly addressSettingsService: AddressSettingsService,
|
||||
) {}
|
||||
|
||||
public async onModuleInit(): Promise<void> {
|
||||
if (process.env.MASTER === 'true') {
|
||||
return;
|
||||
}
|
||||
|
||||
const ports = this.getPorts();
|
||||
if (ports.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
await this.ensureInitialized();
|
||||
ports.forEach(port => this.startSocketServer(port));
|
||||
}
|
||||
|
||||
public async ensureInitialized(): Promise<void> {
|
||||
if (this.noiseConfig != null) {
|
||||
return;
|
||||
}
|
||||
|
||||
await this.initializeNoiseConfig();
|
||||
}
|
||||
|
||||
public createClient(socket: Socket, firstChunk: Buffer): StratumV2Client {
|
||||
if (this.noiseConfig == null) {
|
||||
throw new Error('Stratum V2 service is not initialized');
|
||||
}
|
||||
|
||||
return new StratumV2Client(
|
||||
socket,
|
||||
firstChunk,
|
||||
this,
|
||||
this.stratumV1JobsService,
|
||||
this.bitcoinRpcService,
|
||||
this.clientService,
|
||||
this.clientStatisticsService,
|
||||
this.notificationService,
|
||||
this.blocksService,
|
||||
this.configService,
|
||||
this.addressSettingsService,
|
||||
);
|
||||
}
|
||||
|
||||
public getNoiseConfig(): Sv2NoiseConfig {
|
||||
return this.noiseConfig;
|
||||
}
|
||||
|
||||
public async getPoolAuthorityPublicKey(): Promise<{ publicKey: string; configured: boolean }> {
|
||||
await this.ensureInitialized();
|
||||
|
||||
return {
|
||||
publicKey: encodeSv2AuthorityPublicKey(this.authorityPublicKeyXOnly),
|
||||
configured: this.authorityKeyConfigured,
|
||||
};
|
||||
}
|
||||
|
||||
public getNextChannelId(): number {
|
||||
return this.channelIdCounter++;
|
||||
}
|
||||
|
||||
public generateExtranoncePrefix(): Buffer {
|
||||
const prefix = Buffer.alloc(4);
|
||||
prefix.writeUInt16BE(this.channelIdCounter & 0xffff, 0);
|
||||
crypto.randomBytes(2).copy(prefix, 2);
|
||||
return prefix;
|
||||
}
|
||||
|
||||
public allocateExtendedExtranoncePrefix(channelId: number): Buffer {
|
||||
return this.extranonceManager.allocate(channelId);
|
||||
}
|
||||
|
||||
public releaseExtendedExtranoncePrefix(channelId: number): void {
|
||||
this.extranonceManager.release(channelId);
|
||||
}
|
||||
|
||||
public getExtendedMinerExtranonceSize(): number {
|
||||
return this.extranonceManager.minerExtranonceSize;
|
||||
}
|
||||
|
||||
private async initializeNoiseConfig(): Promise<void> {
|
||||
const configuredAuthorityKey = this.configService.get<string>('SV2_AUTHORITY_PRIVKEY');
|
||||
this.authorityKeyConfigured = configuredAuthorityKey?.length === 64;
|
||||
this.authorityPrivKey = configuredAuthorityKey?.length === 64
|
||||
? Buffer.from(configuredAuthorityKey, 'hex')
|
||||
: crypto.randomBytes(32);
|
||||
this.authorityPublicKeyXOnly = xOnlyPubKeyFromPriv(this.authorityPrivKey);
|
||||
|
||||
if (!configuredAuthorityKey) {
|
||||
console.warn('SV2_AUTHORITY_PRIVKEY is not set; generated an ephemeral SV2 authority key');
|
||||
}
|
||||
|
||||
this.serverKeypair = await generateServerKeypair();
|
||||
const now = Math.floor(Date.now() / 1000);
|
||||
this.noiseConfig = {
|
||||
staticKeypair: this.serverKeypair,
|
||||
certificateMessage: createSignatureNoiseMessage(
|
||||
this.authorityPrivKey,
|
||||
xOnlyPubKeyFromPriv(this.serverKeypair.privateKey),
|
||||
now - 3600,
|
||||
now + 86400,
|
||||
),
|
||||
};
|
||||
}
|
||||
|
||||
private getPorts(): number[] {
|
||||
const configuredPorts = this.configService.get<string>('STRATUM_V2_PORTS');
|
||||
if (!configuredPorts?.trim()) {
|
||||
return [];
|
||||
}
|
||||
|
||||
const ports = configuredPorts
|
||||
.split(',')
|
||||
.map(port => parseInt(port.trim(), 10))
|
||||
.filter(port => Number.isInteger(port) && port > 0 && port <= 65535);
|
||||
|
||||
return Array.from(new Set(ports));
|
||||
}
|
||||
|
||||
private startSocketServer(port: number): void {
|
||||
const server = new Server((socket: Socket) => {
|
||||
socket.setTimeout(this.getSocketTimeoutMs());
|
||||
socket.setNoDelay(true);
|
||||
|
||||
let client: StratumV2Client = null;
|
||||
|
||||
const closeSocket = () => {
|
||||
if (client != null) {
|
||||
void client.destroy();
|
||||
}
|
||||
if (!socket.destroyed) {
|
||||
socket.destroy();
|
||||
}
|
||||
};
|
||||
|
||||
socket.once('data', (firstChunk: Buffer) => {
|
||||
client = this.createClient(socket, firstChunk);
|
||||
});
|
||||
|
||||
socket.on('timeout', closeSocket);
|
||||
socket.on('error', (error: NodeJS.ErrnoException) => {
|
||||
if (error.code !== 'ECONNRESET') {
|
||||
console.error(`Stratum V2 socket error: ${error.message}`);
|
||||
}
|
||||
closeSocket();
|
||||
});
|
||||
socket.on('close', () => {
|
||||
if (client != null) {
|
||||
void client.destroy();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
server.on('error', (error) => {
|
||||
console.error(`Stratum V2 server error on port ${port}: ${error.message}`);
|
||||
});
|
||||
|
||||
server.listen(port, () => {
|
||||
console.log(`Stratum V2 server is listening on port ${port}`);
|
||||
});
|
||||
this.servers.push(server);
|
||||
}
|
||||
|
||||
private getSocketTimeoutMs(): number {
|
||||
const configured = parseInt(this.configService.get<string>('STRATUM_V2_SOCKET_TIMEOUT_MS') ?? '', 10);
|
||||
return Number.isFinite(configured) && configured > 0 ? configured : 1000 * 60 * 15;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user