diff --git a/src/services/stratum-v1.service.spec.ts b/src/services/stratum-v1.service.spec.ts index 41c47e1..1031862 100644 --- a/src/services/stratum-v1.service.spec.ts +++ b/src/services/stratum-v1.service.spec.ts @@ -5,10 +5,12 @@ describe('StratumV1Service', () => { const originalStratumPorts = process.env.STRATUM_PORTS; const originalStratumSecure = process.env.STRATUM_SECURE; const originalSecureStratumPorts = process.env.SECURE_STRATUM_PORTS; + const originalBackpressureEnabled = process.env.STRATUM_BACKPRESSURE_ENABLED; let service: StratumV1Service; let clientService; let consoleLogSpy: jest.SpyInstance; + let consoleWarnSpy: jest.SpyInstance; beforeEach(() => { jest.useFakeTimers(); @@ -26,6 +28,7 @@ describe('StratumV1Service', () => { {} as any ); consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined); + consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined); }); afterEach(() => { @@ -33,7 +36,9 @@ describe('StratumV1Service', () => { restoreEnv('STRATUM_PORTS', originalStratumPorts); restoreEnv('STRATUM_SECURE', originalStratumSecure); restoreEnv('SECURE_STRATUM_PORTS', originalSecureStratumPorts); + restoreEnv('STRATUM_BACKPRESSURE_ENABLED', originalBackpressureEnabled); consoleLogSpy.mockRestore(); + consoleWarnSpy.mockRestore(); jest.useRealTimers(); }); @@ -68,6 +73,52 @@ describe('StratumV1Service', () => { 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) { if (value == null) { delete process.env[key]; diff --git a/src/services/stratum-v1.service.ts b/src/services/stratum-v1.service.ts index 442c6ae..4ba0282 100644 --- a/src/services/stratum-v1.service.ts +++ b/src/services/stratum-v1.service.ts @@ -1,6 +1,7 @@ import { Injectable, OnModuleInit } from '@nestjs/common'; import { ConfigService } from '@nestjs/config'; import { Server, Socket } from 'net'; +import { monitorEventLoopDelay } from 'perf_hooks'; import { StratumV1Client } from '../models/StratumV1Client'; import { AddressSettingsService } from '../ORM/address-settings/address-settings.service'; @@ -15,6 +16,20 @@ import { readFileSync } from 'fs'; import { TlsOptions, TLSSocket, createServer } from 'tls'; 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() @@ -24,6 +39,10 @@ export class StratumV1Service implements OnModuleInit { private emptySocket = 0; private normalClosure = 0; private errorClosure = 0; + private readonly listeners: StratumListenerState[] = []; + private readonly eventLoopDelay = monitorEventLoopDelay({ resolution: 20 }); + private backpressureMonitor: NodeJS.Timeout | null = null; + private healthyBackpressureChecks = 0; constructor( private readonly bitcoinRpcService: BitcoinRpcService, @@ -66,9 +85,22 @@ export class StratumV1Service implements OnModuleInit { this.errorClosure = 0; }, 1000 * 60); + this.startBackpressureMonitor(); + } 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) => { // Set 15-minute timeout socket.setTimeout(1000 * 60 * 15); @@ -131,13 +163,21 @@ export class StratumV1Service implements OnModuleInit { console.error(`Server error: ${err.message}`); }); - server.listen(port, () => { - console.log(`Stratum server is listening on port ${port}`); - }); - + return server; } 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 keyPath = path.join(currentDirectory, 'secrets', 'key.pem'); @@ -203,9 +243,140 @@ export class StratumV1Service implements OnModuleInit { console.error(`Server error: ${err.message}`); }); - server.listen(port, () => { - console.log(`Stratum TLS server is listening on port ${port}`); + return server; + + } + + 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; }