Revert "optimize stratum broadcast fanout"

This reverts commit d9fa933825.
This commit is contained in:
Ben
2026-05-06 21:37:43 -04:00
parent 703546719c
commit 8584ab7efb
4 changed files with 55 additions and 453 deletions
-26
View File
@@ -123,32 +123,6 @@ describe('MiningJob', () => {
expect(clone.ins[0].witness[0]).toHaveLength(32); expect(clone.ins[0].witness[0]).toHaveLength(32);
}); });
it('should clone cached job data with a new job id', () => {
const clone = job.cloneForJobId('2');
const originalNotify = JSON.parse(job.response(jobTemplate));
const cloneNotify = JSON.parse(clone.response(jobTemplate));
expect(originalNotify.params[0]).toBe('1');
expect(cloneNotify.params[0]).toBe('2');
expect(cloneNotify.params[2]).toBe(originalNotify.params[2]);
expect(cloneNotify.params[3]).toBe(originalNotify.params[3]);
expect(clone.buildHeaderBuffer(
jobTemplate,
parseInt('00002000', 16),
parseInt('ed460d91', 16),
'57a6f098',
'c708000000000000',
parseInt(MockRecording1.TIME, 16)
).equals(job.buildHeaderBuffer(
jobTemplate,
parseInt('00002000', 16),
parseInt('ed460d91', 16),
'57a6f098',
'c708000000000000',
parseInt(MockRecording1.TIME, 16)
))).toBe(true);
});
it('should omit oversized pool identifiers from the coinbase script', () => { it('should omit oversized pool identifiers from the coinbase script', () => {
jest.spyOn(console, 'warn').mockImplementation(() => undefined); jest.spyOn(console, 'warn').mockImplementation(() => undefined);
const oversizedIdentifier = 'x'.repeat(120); const oversizedIdentifier = 'x'.repeat(120);
-18
View File
@@ -107,24 +107,6 @@ export class MiningJob {
return this.coinbaseTransaction.__toBuffer().toString('hex'); return this.coinbaseTransaction.__toBuffer().toString('hex');
} }
public cloneForJobId(jobId: string): MiningJob {
const job = Object.create(MiningJob.prototype) as MiningJob;
job.network = this.network;
job.jobId = jobId;
job.coinbaseTransaction = this.coinbaseTransaction;
job.coinbasePart1 = this.coinbasePart1;
job.coinbasePart2 = this.coinbasePart2;
job.coinbasePart1Buffer = this.coinbasePart1Buffer;
job.coinbasePart2Buffer = this.coinbasePart2Buffer;
job.merkleBranchBuffers = this.merkleBranchBuffers;
job.notifyStaticParams = this.notifyStaticParams;
job.jobTemplateId = this.jobTemplateId;
job.networkDifficulty = this.networkDifficulty;
job.creation = new Date().getTime();
job.retiredAt = undefined;
return job;
}
public getCoinbasePrefixBuffer(): Buffer { public getCoinbasePrefixBuffer(): Buffer {
return Buffer.from(this.coinbasePart1Buffer); return Buffer.from(this.coinbasePart1Buffer);
} }
+5 -161
View File
@@ -104,8 +104,6 @@ describe('StratumV1Client', () => {
jest.useFakeTimers({ advanceTimers: true }) jest.useFakeTimers({ advanceTimers: true })
jest.setSystemTime(new Date(parseInt(MockRecording1.TIME, 16) * 1000)); jest.setSystemTime(new Date(parseInt(MockRecording1.TIME, 16) * 1000));
delete process.env.STRATUM_DIFFICULTY_CHECK_INTERVAL_MS;
delete process.env.STRATUM_DIFFICULTY_CHECK_BATCH_SIZE;
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined); consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
consoleErrorSpy = jest.spyOn(console, 'error').mockImplementation(() => undefined); consoleErrorSpy = jest.spyOn(console, 'error').mockImplementation(() => undefined);
consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined); consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined);
@@ -133,24 +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).activeClients.clear(); (StratumV1Client as any).jobBroadcastQueue = [];
(StratumV1Client as any).jobBroadcastSubscription?.unsubscribe();
(StratumV1Client as any).jobBroadcastSubscription = null;
(StratumV1Client as any).latestJobTemplate = null;
(StratumV1Client as any).pendingJobTemplate = null;
(StratumV1Client as any).jobBroadcastIterator = null;
(StratumV1Client as any).jobBroadcastDraining = false; (StratumV1Client as any).jobBroadcastDraining = false;
(StratumV1Client as any).clientJobQueue = []; (StratumV1Client as any).jobBroadcastDrainOffset = 0;
(StratumV1Client as any).clientJobDraining = false;
(StratumV1Client as any).clientJobDrainOffset = 0;
if ((StratumV1Client as any).difficultyCheckInterval != null) {
clearInterval((StratumV1Client as any).difficultyCheckInterval);
}
(StratumV1Client as any).difficultyCheckInterval = null;
(StratumV1Client as any).difficultyCheckIterator = null;
(StratumV1Client as any).difficultyCheckDraining = false;
(StratumV1Client as any).miningJobCacheTemplateId = null;
(StratumV1Client as any).miningJobCache.clear();
bitcoinRpcService = { bitcoinRpcService = {
newBlockTemplate$: newBlockEmitter.asObservable(), newBlockTemplate$: newBlockEmitter.asObservable(),
@@ -449,10 +432,11 @@ describe('StratumV1Client', () => {
}); });
it('should broadcast mining jobs through the shared active-client broadcaster', async () => { it('should queue mining job broadcasts instead of building jobs in the observable callback', async () => {
jest.useFakeTimers(); jest.useFakeTimers();
jest.setSystemTime(new Date(parseInt(MockRecording1.TIME, 16) * 1000)); jest.setSystemTime(new Date(parseInt(MockRecording1.TIME, 16) * 1000));
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); 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'); const sendNewMiningJobSpy = jest.spyOn(client as any, 'sendNewMiningJob');
(client as any).clientSubscription = { userAgent: 'bitaxe' }; (client as any).clientSubscription = { userAgent: 'bitaxe' };
@@ -465,8 +449,7 @@ describe('StratumV1Client', () => {
await Promise.resolve(); await Promise.resolve();
await Promise.resolve(); await Promise.resolve();
expect((StratumV1Client as any).activeClients.has(client)).toBe(true); expect(enqueueNewMiningJobSpy).toHaveBeenCalled();
expect((StratumV1Client as any).jobBroadcastSubscription).not.toBeNull();
jest.advanceTimersByTime(0); jest.advanceTimersByTime(0);
await Promise.resolve(); await Promise.resolve();
@@ -475,145 +458,6 @@ describe('StratumV1Client', () => {
expect(sendNewMiningJobSpy).toHaveBeenCalled(); expect(sendNewMiningJobSpy).toHaveBeenCalled();
}); });
it('should send the current job only to newly initialized clients after the shared broadcaster is running', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
const firstClientSendSpy = 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();
jest.advanceTimersByTime(0);
await Promise.resolve();
await Promise.resolve();
firstClientSendSpy.mockClear();
const secondSocket = new Socket();
jest.spyOn(secondSocket, 'on').mockImplementation(() => secondSocket);
secondSocket.end = jest.fn();
jest.spyOn(secondSocket, 'destroy').mockImplementation(() => secondSocket);
const secondClient = new StratumV1Client(
secondSocket,
stratumV1JobsService,
bitcoinRpcService,
clientService,
clientStatisticsService,
notificationService,
blocksService,
configService,
moduleRef.get<AddressSettingsService>(AddressSettingsService)
);
jest.spyOn(secondClient as any, 'write').mockImplementation((data) => Promise.resolve(true));
const secondClientSendSpy = jest.spyOn(secondClient as any, 'sendNewMiningJob');
(secondClient as any).clientSubscription = { userAgent: 'bitaxe' };
(secondClient as any).clientAuthorization = {
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
worker: 'bitaxe4'
};
(secondClient as any).statistics = {};
await (secondClient as any).initStratum();
jest.advanceTimersByTime(0);
await Promise.resolve();
await Promise.resolve();
expect(secondClientSendSpy).toHaveBeenCalled();
expect(firstClientSendSpy).not.toHaveBeenCalled();
expect((StratumV1Client as any).miningJobCache.size).toBe(1);
await secondClient.destroy();
});
it('should skip queued current-job sends when a newer template arrives first', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
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();
jest.advanceTimersByTime(0);
await Promise.resolve();
await Promise.resolve();
sendNewMiningJobSpy.mockClear();
const oldTemplate = (StratumV1Client as any).latestJobTemplate;
(StratumV1Client as any).queueClientJob(client, oldTemplate);
(StratumV1Client as any).latestJobTemplate = {
...oldTemplate,
blockData: {
...oldTemplate.blockData,
id: 'new-template'
}
};
jest.advanceTimersByTime(0);
await Promise.resolve();
await Promise.resolve();
expect(sendNewMiningJobSpy).not.toHaveBeenCalled();
});
it('should clear duplicate share tracking immediately when a clear-jobs template is queued', async () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
(client as any).clientSubscription = { userAgent: 'bitaxe' };
(client as any).clientAuthorization = {
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
worker: 'bitaxe3'
};
(client as any).statistics = {};
await (client as any).initStratum();
jest.advanceTimersByTime(0);
await Promise.resolve();
await Promise.resolve();
(client as any).miningSubmissionHashes.add('old-share');
const latestTemplate = (StratumV1Client as any).latestJobTemplate;
(StratumV1Client as any).queueJobBroadcast({
...latestTemplate,
blockData: {
...latestTemplate.blockData,
id: 'clear-template',
clearJobs: true
}
});
expect((client as any).miningSubmissionHashes.size).toBe(0);
});
it('should check difficulty through the shared scheduler', async () => {
process.env.STRATUM_DIFFICULTY_CHECK_INTERVAL_MS = '1000';
process.env.STRATUM_DIFFICULTY_CHECK_BATCH_SIZE = '1';
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
const checkDifficultySpy = jest.spyOn(client as any, 'checkDifficulty').mockResolvedValue(undefined);
(client as any).clientSubscription = { userAgent: 'bitaxe' };
(client as any).clientAuthorization = {
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
worker: 'bitaxe3'
};
(client as any).statistics = {};
await (client as any).initStratum();
expect((StratumV1Client as any).activeClients.has(client)).toBe(true);
expect((StratumV1Client as any).difficultyCheckInterval).not.toBeNull();
jest.advanceTimersByTime(1000);
jest.runOnlyPendingTimers();
await Promise.resolve();
await Promise.resolve();
expect(checkDifficultySpy).toHaveBeenCalled();
delete process.env.STRATUM_DIFFICULTY_CHECK_INTERVAL_MS;
delete process.env.STRATUM_DIFFICULTY_CHECK_BATCH_SIZE;
});
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');
+50 -248
View File
@@ -5,6 +5,7 @@ import { validate, ValidationError, ValidatorOptions } from 'class-validator';
import * as crypto from 'crypto'; import * as crypto from 'crypto';
import { Socket } from 'net'; import { Socket } from 'net';
import { firstValueFrom, Subscription } from 'rxjs'; import { firstValueFrom, Subscription } from 'rxjs';
import { clearInterval } from 'timers';
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service'; import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
import { BlocksService } from '../ORM/blocks/blocks.service'; import { BlocksService } from '../ORM/blocks/blocks.service';
@@ -32,20 +33,13 @@ 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_JOB_BROADCAST_BATCH_SIZE = 250;
const DEFAULT_HANDSHAKE_TIMEOUT_MS = 30000; const DEFAULT_HANDSHAKE_TIMEOUT_MS = 30000;
const DEFAULT_DIFFICULTY_CHECK_INTERVAL_MS = 60 * 1000;
const DEFAULT_DIFFICULTY_CHECK_BATCH_SIZE = 1000;
interface ClientJobQueueItem { interface JobBroadcastQueueItem {
client: StratumV1Client; client: StratumV1Client;
jobTemplate: IJobTemplate; jobTemplate: IJobTemplate;
jobTemplateId: string; jobTemplateId: string;
} }
interface PayoutInformation {
address: string;
percent: number;
}
export function effectiveJobDifficulty( export function effectiveJobDifficulty(
jobIdInt: number, jobIdInt: number,
currentDiff: number, currentDiff: number,
@@ -65,25 +59,16 @@ 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 activeClients = new Set<StratumV1Client>(); private static jobBroadcastQueue: JobBroadcastQueueItem[] = [];
private static jobBroadcastSubscription: Subscription | null = null;
private static latestJobTemplate: IJobTemplate | null = null;
private static pendingJobTemplate: IJobTemplate | null = null;
private static jobBroadcastIterator: SetIterator<StratumV1Client> | null = null;
private static jobBroadcastDraining = false; private static jobBroadcastDraining = false;
private static clientJobQueue: ClientJobQueueItem[] = []; private static jobBroadcastDrainOffset = 0;
private static clientJobDraining = false;
private static clientJobDrainOffset = 0;
private static difficultyCheckInterval: NodeJS.Timeout | null = null;
private static difficultyCheckIterator: SetIterator<StratumV1Client> | null = null;
private static difficultyCheckDraining = false;
private static miningJobCacheTemplateId: string | null = null;
private static miningJobCache = new Map<string, MiningJob>();
public clientSubscription: SubscriptionMessage; public clientSubscription: SubscriptionMessage;
private clientConfiguration: ConfigurationMessage; private clientConfiguration: ConfigurationMessage;
private clientAuthorization: AuthorizationMessage; private clientAuthorization: AuthorizationMessage;
private clientSuggestedDifficulty: SuggestDifficulty; private clientSuggestedDifficulty: SuggestDifficulty;
private stratumSubscription: Subscription;
private backgroundWork: NodeJS.Timeout[] = [];
private statistics: StratumV1ClientStatistics; private statistics: StratumV1ClientStatistics;
private stratumInitialized = false; private stratumInitialized = false;
@@ -157,17 +142,22 @@ export class StratumV1Client {
} }
this.destroyed = true; this.destroyed = true;
StratumV1Client.activeClients.delete(this);
StratumV1Client.stopSharedWorkersIfIdle();
if (this.clientEntity?.id) { if (this.clientEntity?.id) {
await this.clientService.delete(this.clientEntity.id); await this.clientService.delete(this.clientEntity.id);
} }
if (this.stratumSubscription != null) {
this.stratumSubscription.unsubscribe();
}
if (this.handshakeTimeout != null) { if (this.handshakeTimeout != null) {
clearTimeout(this.handshakeTimeout); clearTimeout(this.handshakeTimeout);
this.handshakeTimeout = null; this.handshakeTimeout = null;
} }
this.backgroundWork.forEach(work => {
clearInterval(work);
});
} }
private getRandomHexString() { private getRandomHexString() {
@@ -453,13 +443,23 @@ export class StratumV1Client {
} }
} }
const sharedJobBroadcastAlreadyRunning = StratumV1Client.jobBroadcastSubscription != null; this.stratumSubscription = this.stratumV1JobsService.newMiningJob$.subscribe(async (jobTemplate) => {
StratumV1Client.activeClients.add(this); try {
StratumV1Client.ensureSharedJobBroadcast(this.stratumV1JobsService); if(jobTemplate.blockData.clearJobs){
if (sharedJobBroadcastAlreadyRunning && StratumV1Client.latestJobTemplate != null) { this.miningSubmissionHashes.clear();
StratumV1Client.queueClientJob(this, StratumV1Client.latestJobTemplate); }
} this.enqueueNewMiningJob(jobTemplate);
StratumV1Client.ensureSharedDifficultyChecks(); } catch (e) {
await this.socket.end();
console.error(e);
}
});
this.backgroundWork.push(
setInterval(async () => {
await this.checkDifficulty();
}, 60 * 1000)
);
// this.backgroundWork.push( // this.backgroundWork.push(
// setInterval(async () => { // setInterval(async () => {
@@ -468,54 +468,21 @@ export class StratumV1Client {
// ); // );
} }
private static ensureSharedJobBroadcast(stratumV1JobsService: StratumV1JobsService) { private enqueueNewMiningJob(jobTemplate: IJobTemplate) {
if (StratumV1Client.jobBroadcastSubscription != null) { if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) {
return; return;
} }
StratumV1Client.jobBroadcastSubscription = stratumV1JobsService.newMiningJob$.subscribe((jobTemplate) => { this.latestQueuedJobTemplateId = jobTemplate.blockData.id;
StratumV1Client.latestJobTemplate = jobTemplate; StratumV1Client.jobBroadcastQueue.push({
StratumV1Client.queueJobBroadcast(jobTemplate); client: this,
jobTemplate,
jobTemplateId: jobTemplate.blockData.id
}); });
StratumV1Client.scheduleJobBroadcastDrain();
} }
private static stopSharedWorkersIfIdle() { private static scheduleJobBroadcastDrain() {
if (StratumV1Client.activeClients.size > 0) {
return;
}
if (StratumV1Client.jobBroadcastSubscription != null) {
StratumV1Client.jobBroadcastSubscription.unsubscribe();
StratumV1Client.jobBroadcastSubscription = null;
}
if (StratumV1Client.difficultyCheckInterval != null) {
clearInterval(StratumV1Client.difficultyCheckInterval);
StratumV1Client.difficultyCheckInterval = null;
}
StratumV1Client.pendingJobTemplate = null;
StratumV1Client.jobBroadcastIterator = null;
StratumV1Client.jobBroadcastDraining = false;
StratumV1Client.latestJobTemplate = null;
StratumV1Client.clientJobQueue = [];
StratumV1Client.clientJobDraining = false;
StratumV1Client.clientJobDrainOffset = 0;
StratumV1Client.difficultyCheckIterator = null;
StratumV1Client.difficultyCheckDraining = false;
StratumV1Client.miningJobCacheTemplateId = null;
StratumV1Client.miningJobCache.clear();
}
private static queueJobBroadcast(jobTemplate: IJobTemplate) {
StratumV1Client.pendingJobTemplate = jobTemplate;
StratumV1Client.jobBroadcastIterator = StratumV1Client.activeClients.values();
if (jobTemplate.blockData.clearJobs) {
for (const client of StratumV1Client.activeClients) {
client.miningSubmissionHashes.clear();
}
}
if (StratumV1Client.jobBroadcastDraining) { if (StratumV1Client.jobBroadcastDraining) {
return; return;
} }
@@ -524,92 +491,34 @@ export class StratumV1Client {
setTimeout(() => StratumV1Client.drainJobBroadcastQueue(), 0); setTimeout(() => StratumV1Client.drainJobBroadcastQueue(), 0);
} }
private static queueClientJob(client: StratumV1Client, jobTemplate: IJobTemplate) { private static drainJobBroadcastQueue() {
if (client.connectionClosed || client.socket.destroyed || client.socket.writableEnded) {
return;
}
client.latestQueuedJobTemplateId = jobTemplate.blockData.id;
StratumV1Client.clientJobQueue.push({
client,
jobTemplate,
jobTemplateId: jobTemplate.blockData.id
});
if (StratumV1Client.clientJobDraining) {
return;
}
StratumV1Client.clientJobDraining = true;
setTimeout(() => StratumV1Client.drainClientJobQueue(), 0);
}
private static drainClientJobQueue() {
const batchSize = StratumV1Client.getJobBroadcastBatchSize(); const batchSize = StratumV1Client.getJobBroadcastBatchSize();
let processed = 0; let processed = 0;
while (processed < batchSize && StratumV1Client.clientJobDrainOffset < StratumV1Client.clientJobQueue.length) { while (processed < batchSize && StratumV1Client.jobBroadcastDrainOffset < StratumV1Client.jobBroadcastQueue.length) {
const item = StratumV1Client.clientJobQueue[StratumV1Client.clientJobDrainOffset++]; const item = StratumV1Client.jobBroadcastQueue[StratumV1Client.jobBroadcastDrainOffset++];
processed++; processed++;
if (item.client.connectionClosed if (item.client.connectionClosed
|| item.client.socket.destroyed || item.client.socket.destroyed
|| item.client.socket.writableEnded || item.client.socket.writableEnded
|| StratumV1Client.latestJobTemplate?.blockData.id !== item.jobTemplateId
|| item.client.latestQueuedJobTemplateId !== item.jobTemplateId) { || item.client.latestQueuedJobTemplateId !== item.jobTemplateId) {
continue; continue;
} }
item.client.sendNewMiningJob(item.jobTemplate).catch((e) => { item.client.sendNewMiningJob(item.jobTemplate).catch(async (e) => {
item.client.closeSocket(); item.client.closeSocket();
console.error(e); console.error(e);
}); });
} }
if (StratumV1Client.clientJobDrainOffset >= StratumV1Client.clientJobQueue.length) { if (StratumV1Client.jobBroadcastDrainOffset >= StratumV1Client.jobBroadcastQueue.length) {
StratumV1Client.clientJobQueue = []; StratumV1Client.jobBroadcastQueue = [];
StratumV1Client.clientJobDrainOffset = 0; StratumV1Client.jobBroadcastDrainOffset = 0;
StratumV1Client.clientJobDraining = false;
return;
}
setTimeout(() => StratumV1Client.drainClientJobQueue(), 0);
}
private static drainJobBroadcastQueue() {
const jobTemplate = StratumV1Client.pendingJobTemplate;
const iterator = StratumV1Client.jobBroadcastIterator;
if (jobTemplate == null || iterator == null) {
StratumV1Client.jobBroadcastDraining = false; StratumV1Client.jobBroadcastDraining = false;
return; return;
} }
const batchSize = StratumV1Client.getJobBroadcastBatchSize();
let processed = 0;
while (processed < batchSize) {
const next = iterator.next();
if (next.done) {
StratumV1Client.jobBroadcastIterator = null;
StratumV1Client.jobBroadcastDraining = false;
return;
}
const client = next.value;
processed++;
if (client.connectionClosed || client.socket.destroyed || client.socket.writableEnded) {
continue;
}
client.latestQueuedJobTemplateId = jobTemplate.blockData.id;
client.sendNewMiningJob(jobTemplate).catch(async (e) => {
client.closeSocket();
console.error(e);
});
}
setTimeout(() => StratumV1Client.drainJobBroadcastQueue(), 0); setTimeout(() => StratumV1Client.drainJobBroadcastQueue(), 0);
} }
@@ -621,79 +530,9 @@ export class StratumV1Client {
return DEFAULT_JOB_BROADCAST_BATCH_SIZE; return DEFAULT_JOB_BROADCAST_BATCH_SIZE;
} }
private static ensureSharedDifficultyChecks() {
if (StratumV1Client.difficultyCheckInterval != null) {
return;
}
StratumV1Client.difficultyCheckInterval = setInterval(() => {
StratumV1Client.scheduleDifficultyCheckDrain();
}, StratumV1Client.getDifficultyCheckIntervalMs());
}
private static scheduleDifficultyCheckDrain() {
StratumV1Client.difficultyCheckIterator = StratumV1Client.activeClients.values();
if (StratumV1Client.difficultyCheckDraining) {
return;
}
StratumV1Client.difficultyCheckDraining = true;
setTimeout(() => StratumV1Client.drainDifficultyChecks(), 0);
}
private static drainDifficultyChecks() {
const iterator = StratumV1Client.difficultyCheckIterator;
if (iterator == null) {
StratumV1Client.difficultyCheckDraining = false;
return;
}
const batchSize = StratumV1Client.getDifficultyCheckBatchSize();
let processed = 0;
while (processed < batchSize) {
const next = iterator.next();
if (next.done) {
StratumV1Client.difficultyCheckIterator = null;
StratumV1Client.difficultyCheckDraining = false;
return;
}
const client = next.value;
processed++;
if (client.connectionClosed || client.socket.destroyed || client.socket.writableEnded) {
continue;
}
client.checkDifficulty().catch((e) => {
client.closeSocket();
console.error(e);
});
}
setTimeout(() => StratumV1Client.drainDifficultyChecks(), 0);
}
private static getDifficultyCheckIntervalMs() {
const configuredInterval = parseInt(process.env.STRATUM_DIFFICULTY_CHECK_INTERVAL_MS, 10);
if (Number.isFinite(configuredInterval) && configuredInterval > 0) {
return configuredInterval;
}
return DEFAULT_DIFFICULTY_CHECK_INTERVAL_MS;
}
private static getDifficultyCheckBatchSize() {
const configuredBatchSize = parseInt(process.env.STRATUM_DIFFICULTY_CHECK_BATCH_SIZE, 10);
if (Number.isFinite(configuredBatchSize) && configuredBatchSize > 0) {
return configuredBatchSize;
}
return DEFAULT_DIFFICULTY_CHECK_BATCH_SIZE;
}
private async sendNewMiningJob(jobTemplate: IJobTemplate) { private async sendNewMiningJob(jobTemplate: IJobTemplate) {
let payoutInformation: PayoutInformation[] = [ let payoutInformation= [
{ address: this.clientAuthorization.address, percent: 100 } { address: this.clientAuthorization.address, percent: 100 }
]; ];
// const devFeeAddress = this.configService.get('DEV_FEE_ADDRESS'); // const devFeeAddress = this.configService.get('DEV_FEE_ADDRESS');
@@ -732,14 +571,12 @@ export class StratumV1Client {
throw new Error('Invalid network configuration'); throw new Error('Invalid network configuration');
} }
const poolIdentifier = this.configService.get('POOL_IDENTIFIER') || 'Public-Pool'; const job = new MiningJob(
const job = StratumV1Client.getCachedMiningJob(
network, network,
this.stratumV1JobsService.getNextId(), this.stratumV1JobsService.getNextId(),
payoutInformation, payoutInformation,
jobTemplate, jobTemplate,
poolIdentifier, this.configService.get('POOL_IDENTIFIER') || 'Public-Pool'
networkConfig
); );
this.stratumV1JobsService.addJob(job); this.stratumV1JobsService.addJob(job);
@@ -755,41 +592,6 @@ export class StratumV1Client {
} }
private static getCachedMiningJob(
network: bitcoinjs.networks.Network,
jobId: string,
payoutInformation: PayoutInformation[],
jobTemplate: IJobTemplate,
poolIdentifier: string,
networkConfig: string,
) {
if (StratumV1Client.miningJobCacheTemplateId !== jobTemplate.blockData.id) {
StratumV1Client.miningJobCacheTemplateId = jobTemplate.blockData.id;
StratumV1Client.miningJobCache.clear();
}
const cacheKey = JSON.stringify({
template: jobTemplate.blockData.id,
network: networkConfig,
poolIdentifier,
payoutInformation
});
let cachedJob = StratumV1Client.miningJobCache.get(cacheKey);
if (cachedJob == null) {
cachedJob = new MiningJob(
network,
jobId,
payoutInformation,
jobTemplate,
poolIdentifier
);
StratumV1Client.miningJobCache.set(cacheKey, cachedJob);
return cachedJob;
}
return cachedJob.cloneForJobId(jobId);
}
private async ensureClientEntity() { private async ensureClientEntity() {
if (this.clientEntity != null) { if (this.clientEntity != null) {