improve stratum connection and broadcast performance

This commit is contained in:
Ben
2026-05-06 20:38:42 -04:00
parent 8709de3d83
commit d6d3727119
3 changed files with 198 additions and 5 deletions
+64
View File
@@ -131,6 +131,9 @@ describe('StratumV1Client', () => {
}); });
(StratumV1Client as any).blockedUserAgentLogState.clear(); (StratumV1Client as any).blockedUserAgentLogState.clear();
(StratumV1Client as any).validationErrorLogState.clear(); (StratumV1Client as any).validationErrorLogState.clear();
(StratumV1Client as any).jobBroadcastQueue = [];
(StratumV1Client as any).jobBroadcastDraining = false;
(StratumV1Client as any).jobBroadcastDrainOffset = 0;
bitcoinRpcService = { bitcoinRpcService = {
newBlockTemplate$: newBlockEmitter.asObservable(), newBlockTemplate$: newBlockEmitter.asObservable(),
@@ -299,6 +302,41 @@ describe('StratumV1Client', () => {
await secondClient.destroy(); await secondClient.destroy();
}); });
it('should close clients that do not complete the stratum handshake', async () => {
(configService.get as jest.Mock).mockImplementation((key: string) => {
switch (key) {
case 'STRATUM_HANDSHAKE_TIMEOUT_MS':
return '1000';
case 'DEV_FEE_ADDRESS':
return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
case 'NETWORK':
return 'testnet';
}
return null;
});
const slowSocket = new Socket();
jest.spyOn(slowSocket, 'on').mockImplementation(() => slowSocket);
jest.spyOn(slowSocket, 'destroy').mockImplementation(() => slowSocket);
slowSocket.end = jest.fn();
const slowClient = new StratumV1Client(
slowSocket,
stratumV1JobsService,
bitcoinRpcService,
clientService,
clientStatisticsService,
notificationService,
blocksService,
configService,
moduleRef.get<AddressSettingsService>(AddressSettingsService)
);
jest.advanceTimersByTime(1000);
expect(slowSocket.destroy).toHaveBeenCalled();
await slowClient.destroy();
});
it('should respond to mining.configure', async () => { it('should respond to mining.configure', async () => {
@@ -394,6 +432,32 @@ describe('StratumV1Client', () => {
}); });
it('should queue mining job broadcasts instead of building jobs in the observable callback', async () => {
jest.useFakeTimers();
jest.setSystemTime(new Date(parseInt(MockRecording1.TIME, 16) * 1000));
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
const enqueueNewMiningJobSpy = jest.spyOn(client as any, 'enqueueNewMiningJob');
const sendNewMiningJobSpy = jest.spyOn(client as any, 'sendNewMiningJob');
(client as any).clientSubscription = { userAgent: 'bitaxe' };
(client as any).clientAuthorization = {
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
worker: 'bitaxe3'
};
(client as any).statistics = {};
await (client as any).initStratum();
await Promise.resolve();
await Promise.resolve();
expect(enqueueNewMiningJobSpy).toHaveBeenCalled();
jest.advanceTimersByTime(0);
await Promise.resolve();
await Promise.resolve();
expect(sendNewMiningJobSpy).toHaveBeenCalled();
});
it('should use the header-only fast path for non-block submissions', async () => { it('should use the header-only fast path for non-block submissions', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
const buildHeaderSpy = jest.spyOn(MiningJob.prototype, 'buildHeaderBuffer'); const buildHeaderSpy = jest.spyOn(MiningJob.prototype, 'buildHeaderBuffer');
+110 -1
View File
@@ -31,6 +31,14 @@ import { StratumV1ClientStatistics } from './StratumV1ClientStatistics';
const TRUE_DIFF_ONE = 2.695953529101131e67; const TRUE_DIFF_ONE = 2.695953529101131e67;
const BLOCKED_USER_AGENT_LOG_INTERVAL_MS = 60 * 1000; const BLOCKED_USER_AGENT_LOG_INTERVAL_MS = 60 * 1000;
const VALIDATION_ERROR_LOG_INTERVAL_MS = 60 * 1000; const VALIDATION_ERROR_LOG_INTERVAL_MS = 60 * 1000;
const DEFAULT_JOB_BROADCAST_BATCH_SIZE = 250;
const DEFAULT_HANDSHAKE_TIMEOUT_MS = 30000;
interface JobBroadcastQueueItem {
client: StratumV1Client;
jobTemplate: IJobTemplate;
jobTemplateId: string;
}
export function effectiveJobDifficulty( export function effectiveJobDifficulty(
jobIdInt: number, jobIdInt: number,
@@ -51,6 +59,9 @@ export function effectiveJobDifficulty(
export class StratumV1Client { export class StratumV1Client {
private static blockedUserAgentLogState = new Map<string, { nextLogAt: number, suppressed: number }>(); private static blockedUserAgentLogState = new Map<string, { nextLogAt: number, suppressed: number }>();
private static validationErrorLogState = new Map<string, { nextLogAt: number, suppressed: number, sample: string }>(); private static validationErrorLogState = new Map<string, { nextLogAt: number, suppressed: number, sample: string }>();
private static jobBroadcastQueue: JobBroadcastQueueItem[] = [];
private static jobBroadcastDraining = false;
private static jobBroadcastDrainOffset = 0;
public clientSubscription: SubscriptionMessage; public clientSubscription: SubscriptionMessage;
private clientConfiguration: ConfigurationMessage; private clientConfiguration: ConfigurationMessage;
@@ -76,6 +87,9 @@ export class StratumV1Client {
private buffer: string = ''; private buffer: string = '';
private connectionClosed = false; private connectionClosed = false;
private destroyed = false;
private handshakeTimeout: NodeJS.Timeout | null;
private latestQueuedJobTemplateId: string;
private miningSubmissionHashes = new Set<string>() private miningSubmissionHashes = new Set<string>()
@@ -91,6 +105,12 @@ export class StratumV1Client {
private readonly addressSettingsService: AddressSettingsService private readonly addressSettingsService: AddressSettingsService
) { ) {
this.handshakeTimeout = setTimeout(() => {
if (!this.stratumInitialized) {
this.closeSocket();
}
}, this.getHandshakeTimeoutMs());
this.socket.on('data', (data: Buffer) => { this.socket.on('data', (data: Buffer) => {
this.buffer += data.toString(); this.buffer += data.toString();
let lines = this.buffer.split('\n'); let lines = this.buffer.split('\n');
@@ -115,6 +135,12 @@ export class StratumV1Client {
} }
public async destroy() { public async destroy() {
this.connectionClosed = true;
if (this.destroyed) {
return;
}
this.destroyed = true;
if (this.clientEntity?.id) { if (this.clientEntity?.id) {
await this.clientService.delete(this.clientEntity.id); await this.clientService.delete(this.clientEntity.id);
@@ -124,6 +150,11 @@ export class StratumV1Client {
this.stratumSubscription.unsubscribe(); this.stratumSubscription.unsubscribe();
} }
if (this.handshakeTimeout != null) {
clearTimeout(this.handshakeTimeout);
this.handshakeTimeout = null;
}
this.backgroundWork.forEach(work => { this.backgroundWork.forEach(work => {
clearInterval(work); clearInterval(work);
}); });
@@ -385,6 +416,10 @@ export class StratumV1Client {
private async initStratum() { private async initStratum() {
this.stratumInitialized = true; this.stratumInitialized = true;
if (this.handshakeTimeout != null) {
clearTimeout(this.handshakeTimeout);
this.handshakeTimeout = null;
}
if (this.isBlockedUserAgent(this.clientSubscription.userAgent)) { if (this.isBlockedUserAgent(this.clientSubscription.userAgent)) {
this.logBlockedUserAgent(this.clientSubscription.userAgent); this.logBlockedUserAgent(this.clientSubscription.userAgent);
@@ -413,7 +448,7 @@ export class StratumV1Client {
if(jobTemplate.blockData.clearJobs){ if(jobTemplate.blockData.clearJobs){
this.miningSubmissionHashes.clear(); this.miningSubmissionHashes.clear();
} }
await this.sendNewMiningJob(jobTemplate); this.enqueueNewMiningJob(jobTemplate);
} catch (e) { } catch (e) {
await this.socket.end(); await this.socket.end();
console.error(e); console.error(e);
@@ -433,6 +468,68 @@ export class StratumV1Client {
// ); // );
} }
private enqueueNewMiningJob(jobTemplate: IJobTemplate) {
if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) {
return;
}
this.latestQueuedJobTemplateId = jobTemplate.blockData.id;
StratumV1Client.jobBroadcastQueue.push({
client: this,
jobTemplate,
jobTemplateId: jobTemplate.blockData.id
});
StratumV1Client.scheduleJobBroadcastDrain();
}
private static scheduleJobBroadcastDrain() {
if (StratumV1Client.jobBroadcastDraining) {
return;
}
StratumV1Client.jobBroadcastDraining = true;
setTimeout(() => StratumV1Client.drainJobBroadcastQueue(), 0);
}
private static drainJobBroadcastQueue() {
const batchSize = StratumV1Client.getJobBroadcastBatchSize();
let processed = 0;
while (processed < batchSize && StratumV1Client.jobBroadcastDrainOffset < StratumV1Client.jobBroadcastQueue.length) {
const item = StratumV1Client.jobBroadcastQueue[StratumV1Client.jobBroadcastDrainOffset++];
processed++;
if (item.client.connectionClosed
|| item.client.socket.destroyed
|| item.client.socket.writableEnded
|| item.client.latestQueuedJobTemplateId !== item.jobTemplateId) {
continue;
}
item.client.sendNewMiningJob(item.jobTemplate).catch(async (e) => {
item.client.closeSocket();
console.error(e);
});
}
if (StratumV1Client.jobBroadcastDrainOffset >= StratumV1Client.jobBroadcastQueue.length) {
StratumV1Client.jobBroadcastQueue = [];
StratumV1Client.jobBroadcastDrainOffset = 0;
StratumV1Client.jobBroadcastDraining = false;
return;
}
setTimeout(() => StratumV1Client.drainJobBroadcastQueue(), 0);
}
private static getJobBroadcastBatchSize() {
const configuredBatchSize = parseInt(process.env.STRATUM_JOB_BROADCAST_BATCH_SIZE, 10);
if (Number.isFinite(configuredBatchSize) && configuredBatchSize > 0) {
return configuredBatchSize;
}
return DEFAULT_JOB_BROADCAST_BATCH_SIZE;
}
private async sendNewMiningJob(jobTemplate: IJobTemplate) { private async sendNewMiningJob(jobTemplate: IJobTemplate) {
let payoutInformation= [ let payoutInformation= [
@@ -814,8 +911,20 @@ export class StratumV1Client {
return ` sample=${values.join(',')}`; return ` sample=${values.join(',')}`;
} }
private getHandshakeTimeoutMs() {
const configuredTimeout = parseInt(this.configService.get<string>('STRATUM_HANDSHAKE_TIMEOUT_MS') ?? process.env.STRATUM_HANDSHAKE_TIMEOUT_MS, 10);
if (Number.isFinite(configuredTimeout) && configuredTimeout > 0) {
return configuredTimeout;
}
return DEFAULT_HANDSHAKE_TIMEOUT_MS;
}
private closeSocket() { private closeSocket() {
this.connectionClosed = true; this.connectionClosed = true;
if (this.handshakeTimeout != null) {
clearTimeout(this.handshakeTimeout);
this.handshakeTimeout = null;
}
if (!this.socket.destroyed) { if (!this.socket.destroyed) {
this.socket.destroy(); this.socket.destroy();
} }
+24 -4
View File
@@ -70,6 +70,8 @@ export class StratumV1Service implements OnModuleInit {
private startSocketServer(port: number) { private startSocketServer(port: number) {
const server = new Server(async (socket: Socket) => { const server = new Server(async (socket: Socket) => {
socket.setNoDelay(true);
socket.setKeepAlive(true, 60000);
// Set 15-minute timeout // Set 15-minute timeout
socket.setTimeout(1000 * 60 * 15); socket.setTimeout(1000 * 60 * 15);
@@ -86,9 +88,17 @@ export class StratumV1Service implements OnModuleInit {
); );
// Unified cleanup function // Unified cleanup function
let cleanedUp = false;
const cleanup = async (reason: string) => { const cleanup = async (reason: string) => {
if (client.extraNonceAndSessionId != null) { if (cleanedUp) {
await client.destroy(); return;
}
cleanedUp = true;
const initializedClient = client.extraNonceAndSessionId != null;
await client.destroy();
if (initializedClient) {
if (reason == 'Error') { if (reason == 'Error') {
this.errorClosure++; this.errorClosure++;
} else { } else {
@@ -149,6 +159,8 @@ export class StratumV1Service implements OnModuleInit {
}; };
const server = createServer(tlsOptions, async (socket: TLSSocket) => { const server = createServer(tlsOptions, async (socket: TLSSocket) => {
socket.setNoDelay(true);
socket.setKeepAlive(true, 60000);
// Set 15-minute timeout // Set 15-minute timeout
socket.setTimeout(1000 * 60 * 15); socket.setTimeout(1000 * 60 * 15);
@@ -164,9 +176,17 @@ export class StratumV1Service implements OnModuleInit {
this.addressSettingsService this.addressSettingsService
); );
let cleanedUp = false;
const cleanup = async (reason: string) => { const cleanup = async (reason: string) => {
if (client.extraNonceAndSessionId != null) { if (cleanedUp) {
await client.destroy(); return;
}
cleanedUp = true;
const initializedClient = client.extraNonceAndSessionId != null;
await client.destroy();
if (initializedClient) {
if (reason === 'Error') { if (reason === 'Error') {
this.errorClosure++; this.errorClosure++;
} else { } else {