mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 09:05:06 -07:00
reduce stratum rejection overhead
This commit is contained in:
@@ -130,6 +130,7 @@ describe('StratumV1Client', () => {
|
||||
return null;
|
||||
});
|
||||
(StratumV1Client as any).blockedUserAgentLogState.clear();
|
||||
(StratumV1Client as any).validationErrorLogState.clear();
|
||||
|
||||
bitcoinRpcService = {
|
||||
newBlockTemplate$: newBlockEmitter.asObservable(),
|
||||
@@ -540,6 +541,49 @@ describe('StratumV1Client', () => {
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
|
||||
expect((client as any).write).lastCalledWith(expect.stringContaining(`"error":[20,"Mining Submit validation error"`));
|
||||
expect(socket.destroy).toHaveBeenCalled();
|
||||
expect(consoleWarnSpy).toHaveBeenCalledWith(expect.stringContaining('Mining Submit validation error: extraNonce2:isLength'));
|
||||
});
|
||||
|
||||
it('should throttle repeated mining submit validation logs', async () => {
|
||||
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
|
||||
|
||||
emitMessage(MockRecording1.MINING_SUBSCRIBE);
|
||||
emitMessage(MockRecording1.MINING_AUTHORIZE);
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
|
||||
emitMessage(`{"id": 5, "method": "mining.submit", "params": ["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.bitaxe3", "1", "c7080000", "64b3f3ec", "ed460d91", "00002000"]}`);
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
|
||||
const secondSocket = new Socket();
|
||||
jest.spyOn(secondSocket, 'on').mockImplementation((event: string, listener: (...args: any[]) => void) => {
|
||||
socketEmitter = listener;
|
||||
return secondSocket;
|
||||
});
|
||||
secondSocket.end = jest.fn();
|
||||
jest.spyOn(secondSocket, 'destroy').mockImplementation(() => secondSocket);
|
||||
const secondClient = new StratumV1Client(
|
||||
secondSocket,
|
||||
stratumV1JobsService,
|
||||
bitcoinRpcService,
|
||||
clientService,
|
||||
clientStatisticsService,
|
||||
notificationService,
|
||||
blocksService,
|
||||
configService,
|
||||
moduleRef.get<AddressSettingsService>(AddressSettingsService)
|
||||
);
|
||||
jest.spyOn(secondClient as any, 'write').mockImplementation((data) => Promise.resolve(true));
|
||||
jest.spyOn(secondClient as any, 'getRandomHexString').mockReturnValue(MockRecording1.EXTRA_NONCE);
|
||||
|
||||
socketEmitter(Buffer.from(`${MockRecording1.MINING_SUBSCRIBE}\n`));
|
||||
socketEmitter(Buffer.from(`${MockRecording1.MINING_AUTHORIZE}\n`));
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
socketEmitter(Buffer.from(`{"id": 5, "method": "mining.submit", "params": ["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.bitaxe3", "1", "c7080000", "64b3f3ec", "ed460d91", "00002000"]}\n`));
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
|
||||
expect(consoleWarnSpy.mock.calls.filter(call => call[0]?.startsWith('Mining Submit validation error'))).toHaveLength(1);
|
||||
await secondClient.destroy();
|
||||
});
|
||||
|
||||
it('should close socket when a submit arrives before stratum is initialized', async () => {
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import { ConfigService } from '@nestjs/config';
|
||||
import * as bitcoinjs from 'bitcoinjs-lib';
|
||||
import { plainToInstance } from 'class-transformer';
|
||||
import { validate, ValidatorOptions } from 'class-validator';
|
||||
import { validate, ValidationError, ValidatorOptions } from 'class-validator';
|
||||
import * as crypto from 'crypto';
|
||||
import { Socket } from 'net';
|
||||
import { firstValueFrom, Subscription } from 'rxjs';
|
||||
@@ -30,6 +30,7 @@ import { StratumV1ClientStatistics } from './StratumV1ClientStatistics';
|
||||
|
||||
const TRUE_DIFF_ONE = 2.695953529101131e67;
|
||||
const BLOCKED_USER_AGENT_LOG_INTERVAL_MS = 60 * 1000;
|
||||
const VALIDATION_ERROR_LOG_INTERVAL_MS = 60 * 1000;
|
||||
|
||||
export function effectiveJobDifficulty(
|
||||
jobIdInt: number,
|
||||
@@ -49,6 +50,7 @@ export function effectiveJobDifficulty(
|
||||
|
||||
export class StratumV1Client {
|
||||
private static blockedUserAgentLogState = new Map<string, { nextLogAt: number, suppressed: number }>();
|
||||
private static validationErrorLogState = new Map<string, { nextLogAt: number, suppressed: number, sample: string }>();
|
||||
|
||||
public clientSubscription: SubscriptionMessage;
|
||||
private clientConfiguration: ConfigurationMessage;
|
||||
@@ -348,17 +350,18 @@ export class StratumV1Client {
|
||||
|
||||
|
||||
} else {
|
||||
console.log('Mining Submit validation error');
|
||||
this.logValidationError('Mining Submit validation error', errors);
|
||||
const err = new StratumErrorMessage(
|
||||
miningSubmitMessage.id,
|
||||
eStratumErrorCode.OtherUnknown,
|
||||
'Mining Submit validation error',
|
||||
errors).response();
|
||||
console.error(err);
|
||||
const success = await this.write(err);
|
||||
if (!success) {
|
||||
return;
|
||||
}
|
||||
this.closeSocket();
|
||||
return;
|
||||
}
|
||||
break;
|
||||
}
|
||||
@@ -765,6 +768,52 @@ export class StratumV1Client {
|
||||
});
|
||||
}
|
||||
|
||||
private logValidationError(label: string, errors: ValidationError[]) {
|
||||
const now = Date.now();
|
||||
const signature = this.getValidationErrorSignature(errors);
|
||||
const sample = this.getValidationErrorSample(errors);
|
||||
const key = `${label}:${signature}`;
|
||||
const logState = StratumV1Client.validationErrorLogState.get(key);
|
||||
|
||||
if (logState != null && now < logState.nextLogAt) {
|
||||
logState.suppressed += 1;
|
||||
return;
|
||||
}
|
||||
|
||||
const suppressed = logState?.suppressed ?? 0;
|
||||
const suffix = suppressed > 0 ? ` (${suppressed} similar validation errors suppressed)` : '';
|
||||
console.warn(`${label}: ${signature}${sample}${suffix}`);
|
||||
StratumV1Client.validationErrorLogState.set(key, {
|
||||
nextLogAt: now + VALIDATION_ERROR_LOG_INTERVAL_MS,
|
||||
suppressed: 0,
|
||||
sample
|
||||
});
|
||||
}
|
||||
|
||||
private getValidationErrorSignature(errors: ValidationError[]): string {
|
||||
if (errors.length === 0) {
|
||||
return 'unknown';
|
||||
}
|
||||
|
||||
return errors.map(error => {
|
||||
const constraints = Object.keys(error.constraints ?? {}).sort().join('|') || 'invalid';
|
||||
return `${error.property}:${constraints}`;
|
||||
}).join(';');
|
||||
}
|
||||
|
||||
private getValidationErrorSample(errors: ValidationError[]): string {
|
||||
const values = errors
|
||||
.map(error => error.value)
|
||||
.filter(value => value != null)
|
||||
.map(value => String(value).replace(/[\r\n]/g, '').slice(0, 64));
|
||||
|
||||
if (values.length === 0) {
|
||||
return '';
|
||||
}
|
||||
|
||||
return ` sample=${values.join(',')}`;
|
||||
}
|
||||
|
||||
private closeSocket() {
|
||||
this.connectionClosed = true;
|
||||
if (!this.socket.destroyed) {
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
import { StratumV1Service } from './stratum-v1.service';
|
||||
|
||||
describe('StratumV1Service', () => {
|
||||
const originalMaster = process.env.MASTER;
|
||||
const originalStratumPorts = process.env.STRATUM_PORTS;
|
||||
const originalStratumSecure = process.env.STRATUM_SECURE;
|
||||
const originalSecureStratumPorts = process.env.SECURE_STRATUM_PORTS;
|
||||
|
||||
let service: StratumV1Service;
|
||||
let clientService;
|
||||
let consoleLogSpy: jest.SpyInstance;
|
||||
|
||||
beforeEach(() => {
|
||||
jest.useFakeTimers();
|
||||
clientService = {
|
||||
deleteAll: jest.fn().mockResolvedValue(undefined)
|
||||
};
|
||||
service = new StratumV1Service(
|
||||
{} as any,
|
||||
clientService,
|
||||
{} as any,
|
||||
{} as any,
|
||||
{} as any,
|
||||
{} as any,
|
||||
{} as any,
|
||||
{} as any
|
||||
);
|
||||
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
restoreEnv('MASTER', originalMaster);
|
||||
restoreEnv('STRATUM_PORTS', originalStratumPorts);
|
||||
restoreEnv('STRATUM_SECURE', originalStratumSecure);
|
||||
restoreEnv('SECURE_STRATUM_PORTS', originalSecureStratumPorts);
|
||||
consoleLogSpy.mockRestore();
|
||||
jest.useRealTimers();
|
||||
});
|
||||
|
||||
it('should skip Stratum listeners in the master process', async () => {
|
||||
process.env.MASTER = 'true';
|
||||
const startSocketServerSpy = jest.spyOn(service as any, 'startSocketServer');
|
||||
const startSecureSocketServerSpy = jest.spyOn(service as any, 'startSecureSocketServer');
|
||||
|
||||
await service.onModuleInit();
|
||||
jest.runOnlyPendingTimers();
|
||||
|
||||
expect(clientService.deleteAll).toHaveBeenCalled();
|
||||
expect(startSocketServerSpy).not.toHaveBeenCalled();
|
||||
expect(startSecureSocketServerSpy).not.toHaveBeenCalled();
|
||||
expect(consoleLogSpy).toHaveBeenCalledWith('Master process skipping Stratum socket listeners');
|
||||
});
|
||||
|
||||
it('should start Stratum listeners in worker processes', async () => {
|
||||
process.env.MASTER = 'false';
|
||||
process.env.STRATUM_PORTS = '3333,3334';
|
||||
process.env.STRATUM_SECURE = 'true';
|
||||
process.env.SECURE_STRATUM_PORTS = '4333';
|
||||
const startSocketServerSpy = jest.spyOn(service as any, 'startSocketServer').mockImplementation(() => undefined);
|
||||
const startSecureSocketServerSpy = jest.spyOn(service as any, 'startSecureSocketServer').mockImplementation(() => undefined);
|
||||
|
||||
await service.onModuleInit();
|
||||
jest.advanceTimersByTime(10000);
|
||||
|
||||
expect(clientService.deleteAll).not.toHaveBeenCalled();
|
||||
expect(startSocketServerSpy).toHaveBeenCalledWith(3333);
|
||||
expect(startSocketServerSpy).toHaveBeenCalledWith(3334);
|
||||
expect(startSecureSocketServerSpy).toHaveBeenCalledWith(4333);
|
||||
});
|
||||
|
||||
function restoreEnv(key: string, value: string | undefined) {
|
||||
if (value == null) {
|
||||
delete process.env[key];
|
||||
return;
|
||||
}
|
||||
process.env[key] = value;
|
||||
}
|
||||
});
|
||||
@@ -42,6 +42,8 @@ export class StratumV1Service implements OnModuleInit {
|
||||
|
||||
if (process.env.MASTER == 'true') {
|
||||
await this.clientService.deleteAll();
|
||||
console.log('Master process skipping Stratum socket listeners');
|
||||
return;
|
||||
}
|
||||
|
||||
// wait for all the other processes to init for an even connection distribution
|
||||
@@ -208,4 +210,4 @@ export class StratumV1Service implements OnModuleInit {
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user