mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 09:05:06 -07:00
add stratum worker backpressure
This commit is contained in:
@@ -5,10 +5,12 @@ describe('StratumV1Service', () => {
|
|||||||
const originalStratumPorts = process.env.STRATUM_PORTS;
|
const originalStratumPorts = process.env.STRATUM_PORTS;
|
||||||
const originalStratumSecure = process.env.STRATUM_SECURE;
|
const originalStratumSecure = process.env.STRATUM_SECURE;
|
||||||
const originalSecureStratumPorts = process.env.SECURE_STRATUM_PORTS;
|
const originalSecureStratumPorts = process.env.SECURE_STRATUM_PORTS;
|
||||||
|
const originalBackpressureEnabled = process.env.STRATUM_BACKPRESSURE_ENABLED;
|
||||||
|
|
||||||
let service: StratumV1Service;
|
let service: StratumV1Service;
|
||||||
let clientService;
|
let clientService;
|
||||||
let consoleLogSpy: jest.SpyInstance;
|
let consoleLogSpy: jest.SpyInstance;
|
||||||
|
let consoleWarnSpy: jest.SpyInstance;
|
||||||
|
|
||||||
beforeEach(() => {
|
beforeEach(() => {
|
||||||
jest.useFakeTimers();
|
jest.useFakeTimers();
|
||||||
@@ -26,6 +28,7 @@ describe('StratumV1Service', () => {
|
|||||||
{} as any
|
{} as any
|
||||||
);
|
);
|
||||||
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
|
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
|
||||||
|
consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined);
|
||||||
});
|
});
|
||||||
|
|
||||||
afterEach(() => {
|
afterEach(() => {
|
||||||
@@ -33,7 +36,9 @@ describe('StratumV1Service', () => {
|
|||||||
restoreEnv('STRATUM_PORTS', originalStratumPorts);
|
restoreEnv('STRATUM_PORTS', originalStratumPorts);
|
||||||
restoreEnv('STRATUM_SECURE', originalStratumSecure);
|
restoreEnv('STRATUM_SECURE', originalStratumSecure);
|
||||||
restoreEnv('SECURE_STRATUM_PORTS', originalSecureStratumPorts);
|
restoreEnv('SECURE_STRATUM_PORTS', originalSecureStratumPorts);
|
||||||
|
restoreEnv('STRATUM_BACKPRESSURE_ENABLED', originalBackpressureEnabled);
|
||||||
consoleLogSpy.mockRestore();
|
consoleLogSpy.mockRestore();
|
||||||
|
consoleWarnSpy.mockRestore();
|
||||||
jest.useRealTimers();
|
jest.useRealTimers();
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -68,6 +73,52 @@ describe('StratumV1Service', () => {
|
|||||||
expect(startSecureSocketServerSpy).toHaveBeenCalledWith(4333);
|
expect(startSecureSocketServerSpy).toHaveBeenCalledWith(4333);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('should pause listeners when worker backpressure is high', () => {
|
||||||
|
const close = jest.fn((callback?: (error?: Error) => void) => callback?.());
|
||||||
|
(service as any).listeners.push({
|
||||||
|
port: 3333,
|
||||||
|
secure: false,
|
||||||
|
server: { close },
|
||||||
|
paused: false
|
||||||
|
});
|
||||||
|
jest.spyOn(service as any, 'getEventLoopP95Ms').mockReturnValue(5000);
|
||||||
|
jest.spyOn(service as any, 'getBackpressureEventLoopP95Ms').mockReturnValue(2000);
|
||||||
|
|
||||||
|
(service as any).checkBackpressure();
|
||||||
|
|
||||||
|
expect(close).toHaveBeenCalled();
|
||||||
|
expect((service as any).listeners[0].paused).toBe(true);
|
||||||
|
expect((service as any).listeners[0].server).toBeNull();
|
||||||
|
expect(consoleWarnSpy).toHaveBeenCalledWith(expect.stringContaining('Pausing Stratum accepts'));
|
||||||
|
});
|
||||||
|
|
||||||
|
it('should resume listeners after consecutive healthy backpressure checks', () => {
|
||||||
|
(service as any).listeners.push({
|
||||||
|
port: 3333,
|
||||||
|
secure: false,
|
||||||
|
server: null,
|
||||||
|
paused: true
|
||||||
|
});
|
||||||
|
jest.spyOn(service as any, 'getEventLoopP95Ms').mockReturnValue(50);
|
||||||
|
jest.spyOn(service as any, 'getBackpressureEventLoopP95Ms').mockReturnValue(2000);
|
||||||
|
jest.spyOn(service as any, 'getBackpressureResumeEventLoopP95Ms').mockReturnValue(250);
|
||||||
|
jest.spyOn(service as any, 'getBackpressureResumeRssMb').mockReturnValue(Number.MAX_SAFE_INTEGER);
|
||||||
|
jest.spyOn(service as any, 'getBackpressureHealthyChecks').mockReturnValue(2);
|
||||||
|
const listenSpy = jest.spyOn(service as any, 'listen').mockImplementation((listener: any) => {
|
||||||
|
listener.server = {};
|
||||||
|
listener.paused = false;
|
||||||
|
});
|
||||||
|
|
||||||
|
(service as any).checkBackpressure();
|
||||||
|
expect(listenSpy).not.toHaveBeenCalled();
|
||||||
|
|
||||||
|
(service as any).checkBackpressure();
|
||||||
|
|
||||||
|
expect(listenSpy).toHaveBeenCalledWith((service as any).listeners[0]);
|
||||||
|
expect((service as any).listeners[0].paused).toBe(false);
|
||||||
|
expect(consoleWarnSpy).toHaveBeenCalledWith(expect.stringContaining('Resuming Stratum accepts'));
|
||||||
|
});
|
||||||
|
|
||||||
function restoreEnv(key: string, value: string | undefined) {
|
function restoreEnv(key: string, value: string | undefined) {
|
||||||
if (value == null) {
|
if (value == null) {
|
||||||
delete process.env[key];
|
delete process.env[key];
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import { Injectable, OnModuleInit } from '@nestjs/common';
|
import { Injectable, OnModuleInit } from '@nestjs/common';
|
||||||
import { ConfigService } from '@nestjs/config';
|
import { ConfigService } from '@nestjs/config';
|
||||||
import { Server, Socket } from 'net';
|
import { Server, Socket } from 'net';
|
||||||
|
import { monitorEventLoopDelay } from 'perf_hooks';
|
||||||
|
|
||||||
import { StratumV1Client } from '../models/StratumV1Client';
|
import { StratumV1Client } from '../models/StratumV1Client';
|
||||||
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
|
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
|
||||||
@@ -15,6 +16,20 @@ import { readFileSync } from 'fs';
|
|||||||
import { TlsOptions, TLSSocket, createServer } from 'tls';
|
import { TlsOptions, TLSSocket, createServer } from 'tls';
|
||||||
import * as path from 'path';
|
import * as path from 'path';
|
||||||
|
|
||||||
|
interface StratumListenerState {
|
||||||
|
port: number;
|
||||||
|
secure: boolean;
|
||||||
|
server: Server | null;
|
||||||
|
paused: boolean;
|
||||||
|
}
|
||||||
|
|
||||||
|
const DEFAULT_BACKPRESSURE_CHECK_INTERVAL_MS = 5000;
|
||||||
|
const DEFAULT_BACKPRESSURE_EVENT_LOOP_P95_MS = 2000;
|
||||||
|
const DEFAULT_BACKPRESSURE_EVENT_LOOP_RESUME_P95_MS = 250;
|
||||||
|
const DEFAULT_BACKPRESSURE_RSS_MB = 2500;
|
||||||
|
const DEFAULT_BACKPRESSURE_RESUME_RSS_MB = 2000;
|
||||||
|
const DEFAULT_BACKPRESSURE_HEALTHY_CHECKS = 3;
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
@Injectable()
|
@Injectable()
|
||||||
@@ -24,6 +39,10 @@ export class StratumV1Service implements OnModuleInit {
|
|||||||
private emptySocket = 0;
|
private emptySocket = 0;
|
||||||
private normalClosure = 0;
|
private normalClosure = 0;
|
||||||
private errorClosure = 0;
|
private errorClosure = 0;
|
||||||
|
private readonly listeners: StratumListenerState[] = [];
|
||||||
|
private readonly eventLoopDelay = monitorEventLoopDelay({ resolution: 20 });
|
||||||
|
private backpressureMonitor: NodeJS.Timeout | null = null;
|
||||||
|
private healthyBackpressureChecks = 0;
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
private readonly bitcoinRpcService: BitcoinRpcService,
|
private readonly bitcoinRpcService: BitcoinRpcService,
|
||||||
@@ -66,9 +85,22 @@ export class StratumV1Service implements OnModuleInit {
|
|||||||
this.errorClosure = 0;
|
this.errorClosure = 0;
|
||||||
}, 1000 * 60);
|
}, 1000 * 60);
|
||||||
|
|
||||||
|
this.startBackpressureMonitor();
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private startSocketServer(port: number) {
|
private startSocketServer(port: number) {
|
||||||
|
const listener: StratumListenerState = {
|
||||||
|
port,
|
||||||
|
secure: false,
|
||||||
|
server: null,
|
||||||
|
paused: false
|
||||||
|
};
|
||||||
|
this.listeners.push(listener);
|
||||||
|
this.listen(listener);
|
||||||
|
}
|
||||||
|
|
||||||
|
private createSocketServer(): Server {
|
||||||
const server = new Server(async (socket: Socket) => {
|
const server = new Server(async (socket: Socket) => {
|
||||||
// Set 15-minute timeout
|
// Set 15-minute timeout
|
||||||
socket.setTimeout(1000 * 60 * 15);
|
socket.setTimeout(1000 * 60 * 15);
|
||||||
@@ -131,13 +163,21 @@ export class StratumV1Service implements OnModuleInit {
|
|||||||
console.error(`Server error: ${err.message}`);
|
console.error(`Server error: ${err.message}`);
|
||||||
});
|
});
|
||||||
|
|
||||||
server.listen(port, () => {
|
return server;
|
||||||
console.log(`Stratum server is listening on port ${port}`);
|
|
||||||
});
|
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private startSecureSocketServer(port: number) {
|
private startSecureSocketServer(port: number) {
|
||||||
|
const listener: StratumListenerState = {
|
||||||
|
port,
|
||||||
|
secure: true,
|
||||||
|
server: null,
|
||||||
|
paused: false
|
||||||
|
};
|
||||||
|
this.listeners.push(listener);
|
||||||
|
this.listen(listener);
|
||||||
|
}
|
||||||
|
|
||||||
|
private createSecureSocketServer(): Server {
|
||||||
|
|
||||||
const currentDirectory = process.cwd();
|
const currentDirectory = process.cwd();
|
||||||
const keyPath = path.join(currentDirectory, 'secrets', 'key.pem');
|
const keyPath = path.join(currentDirectory, 'secrets', 'key.pem');
|
||||||
@@ -203,9 +243,140 @@ export class StratumV1Service implements OnModuleInit {
|
|||||||
console.error(`Server error: ${err.message}`);
|
console.error(`Server error: ${err.message}`);
|
||||||
});
|
});
|
||||||
|
|
||||||
server.listen(port, () => {
|
return server;
|
||||||
console.log(`Stratum TLS server is listening on port ${port}`);
|
|
||||||
|
}
|
||||||
|
|
||||||
|
private listen(listener: StratumListenerState) {
|
||||||
|
if (listener.server != null) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
const server = listener.secure ? this.createSecureSocketServer() : this.createSocketServer();
|
||||||
|
listener.server = server;
|
||||||
|
listener.paused = false;
|
||||||
|
|
||||||
|
server.listen(listener.port, () => {
|
||||||
|
console.log(`${listener.secure ? 'Stratum TLS' : 'Stratum'} server is listening on port ${listener.port}`);
|
||||||
});
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
private startBackpressureMonitor() {
|
||||||
|
if (this.isBackpressureDisabled() || this.backpressureMonitor != null) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
this.eventLoopDelay.enable();
|
||||||
|
this.backpressureMonitor = setInterval(() => {
|
||||||
|
this.checkBackpressure();
|
||||||
|
}, this.getBackpressureCheckIntervalMs());
|
||||||
|
}
|
||||||
|
|
||||||
|
private checkBackpressure() {
|
||||||
|
const eventLoopP95Ms = this.getEventLoopP95Ms();
|
||||||
|
const rssMb = Math.round(process.memoryUsage().rss / 1024 / 1024);
|
||||||
|
const overloaded = eventLoopP95Ms >= this.getBackpressureEventLoopP95Ms()
|
||||||
|
|| rssMb >= this.getBackpressureRssMb();
|
||||||
|
const paused = this.listeners.some(listener => listener.paused);
|
||||||
|
|
||||||
|
if (overloaded) {
|
||||||
|
this.healthyBackpressureChecks = 0;
|
||||||
|
if (!paused) {
|
||||||
|
this.pauseAccepting(eventLoopP95Ms, rssMb);
|
||||||
|
}
|
||||||
|
this.eventLoopDelay.reset();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (!paused) {
|
||||||
|
this.eventLoopDelay.reset();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
const healthy = eventLoopP95Ms <= this.getBackpressureResumeEventLoopP95Ms()
|
||||||
|
&& rssMb <= this.getBackpressureResumeRssMb();
|
||||||
|
if (!healthy) {
|
||||||
|
this.healthyBackpressureChecks = 0;
|
||||||
|
this.eventLoopDelay.reset();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
this.healthyBackpressureChecks++;
|
||||||
|
if (this.healthyBackpressureChecks >= this.getBackpressureHealthyChecks()) {
|
||||||
|
this.resumeAccepting(eventLoopP95Ms, rssMb);
|
||||||
|
this.healthyBackpressureChecks = 0;
|
||||||
|
}
|
||||||
|
|
||||||
|
this.eventLoopDelay.reset();
|
||||||
|
}
|
||||||
|
|
||||||
|
private pauseAccepting(eventLoopP95Ms: number, rssMb: number) {
|
||||||
|
console.warn(`Pausing Stratum accepts: eventLoopP95Ms=${eventLoopP95Ms}, rssMb=${rssMb}`);
|
||||||
|
for (const listener of this.listeners) {
|
||||||
|
if (listener.paused || listener.server == null) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
const server = listener.server;
|
||||||
|
listener.server = null;
|
||||||
|
listener.paused = true;
|
||||||
|
server.close((error) => {
|
||||||
|
if (error != null) {
|
||||||
|
console.error(`Error while pausing Stratum listener on port ${listener.port}: ${error.message}`);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private resumeAccepting(eventLoopP95Ms: number, rssMb: number) {
|
||||||
|
console.warn(`Resuming Stratum accepts: eventLoopP95Ms=${eventLoopP95Ms}, rssMb=${rssMb}`);
|
||||||
|
for (const listener of this.listeners) {
|
||||||
|
if (!listener.paused || listener.server != null) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
this.listen(listener);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private getEventLoopP95Ms() {
|
||||||
|
return Math.round(this.eventLoopDelay.percentile(95) / 1e6);
|
||||||
|
}
|
||||||
|
|
||||||
|
private isBackpressureDisabled() {
|
||||||
|
return process.env.STRATUM_BACKPRESSURE_ENABLED?.toLowerCase() === 'false';
|
||||||
|
}
|
||||||
|
|
||||||
|
private getBackpressureCheckIntervalMs() {
|
||||||
|
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_CHECK_INTERVAL_MS', DEFAULT_BACKPRESSURE_CHECK_INTERVAL_MS);
|
||||||
|
}
|
||||||
|
|
||||||
|
private getBackpressureEventLoopP95Ms() {
|
||||||
|
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_EVENT_LOOP_P95_MS', DEFAULT_BACKPRESSURE_EVENT_LOOP_P95_MS);
|
||||||
|
}
|
||||||
|
|
||||||
|
private getBackpressureResumeEventLoopP95Ms() {
|
||||||
|
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_EVENT_LOOP_RESUME_P95_MS', DEFAULT_BACKPRESSURE_EVENT_LOOP_RESUME_P95_MS);
|
||||||
|
}
|
||||||
|
|
||||||
|
private getBackpressureRssMb() {
|
||||||
|
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_RSS_MB', DEFAULT_BACKPRESSURE_RSS_MB);
|
||||||
|
}
|
||||||
|
|
||||||
|
private getBackpressureResumeRssMb() {
|
||||||
|
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_RESUME_RSS_MB', DEFAULT_BACKPRESSURE_RESUME_RSS_MB);
|
||||||
|
}
|
||||||
|
|
||||||
|
private getBackpressureHealthyChecks() {
|
||||||
|
return this.getPositiveIntegerEnv('STRATUM_BACKPRESSURE_HEALTHY_CHECKS', DEFAULT_BACKPRESSURE_HEALTHY_CHECKS);
|
||||||
|
}
|
||||||
|
|
||||||
|
private getPositiveIntegerEnv(key: string, fallback: number) {
|
||||||
|
const configured = parseInt(process.env[key], 10);
|
||||||
|
if (Number.isFinite(configured) && configured > 0) {
|
||||||
|
return configured;
|
||||||
|
}
|
||||||
|
return fallback;
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user