Harden Stratum V1 cleanup

This commit is contained in:
Ben
2026-06-09 12:30:03 -04:00
parent a693501648
commit 9b99c5153a
5 changed files with 139 additions and 46 deletions
+1
View File
@@ -28,6 +28,7 @@ DOCKER_LOG_MAX_FILES=5
# Plain TCP Stratum ports accept both SV1 JSON-RPC and SV2 Noise/binary traffic. # Plain TCP Stratum ports accept both SV1 JSON-RPC and SV2 Noise/binary traffic.
STRATUM_PORTS=3333,3332,3331,3330 STRATUM_PORTS=3333,3332,3331,3330
STRATUM_WORKERS=2 STRATUM_WORKERS=2
STRATUM_WORKER_MAX_MEMORY_RESTART=4096M
STRATUM_MIN_DIFFICULTY=1 STRATUM_MIN_DIFFICULTY=1
STRATUM_SOCKET_TIMEOUT_MS=3600000 STRATUM_SOCKET_TIMEOUT_MS=3600000
STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS=60000 STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS=60000
+1
View File
@@ -42,6 +42,7 @@ module.exports = {
script: './dist/main.js', script: './dist/main.js',
instances: parseInt(process.env.STRATUM_WORKERS || '2', 10), instances: parseInt(process.env.STRATUM_WORKERS || '2', 10),
exec_mode: "cluster", exec_mode: "cluster",
max_memory_restart: process.env.STRATUM_WORKER_MAX_MEMORY_RESTART || '4096M',
env: { env: {
MASTER: 'false', MASTER: 'false',
API_ENABLED: 'false', API_ENABLED: 'false',
+30
View File
@@ -168,6 +168,36 @@ describe('StratumV1Client', () => {
expect(socket.on).toHaveBeenCalled(); expect(socket.on).toHaveBeenCalled();
}); });
it('should clean up socket state only once when destroyed repeatedly', async () => {
const unsubscribe = jest.fn();
const timer = setInterval(() => undefined, 1000);
const removeListenerSpy = jest.spyOn(socket, 'removeListener');
(client as any).clientEntity = {
id: '00000000-0000-4000-8000-000000000001',
address: 'tb1qcleanup',
};
(client as any).stratumSubscription = { unsubscribe };
(client as any).backgroundWork = [timer];
(client as any).miningSubmissionHashes.add('submitted-share');
(client as any).buffer = 'partial-message';
await Promise.all([client.destroy(), client.destroy()]);
expect(redisMessagingService.removeClientPresence).toHaveBeenCalledTimes(1);
expect(redisMessagingService.removeClientPresence).toHaveBeenCalledWith(
'00000000-0000-4000-8000-000000000001',
'tb1qcleanup',
);
expect(clientService.delete).toHaveBeenCalledTimes(1);
expect(clientService.delete).toHaveBeenCalledWith('00000000-0000-4000-8000-000000000001');
expect(unsubscribe).toHaveBeenCalledTimes(1);
expect(removeListenerSpy).toHaveBeenCalledWith('data', expect.any(Function));
expect((client as any).backgroundWork).toEqual([]);
expect((client as any).miningSubmissionHashes.size).toBe(0);
expect((client as any).buffer).toBe('');
});
it('should close socket on invalid JSON', () => { it('should close socket on invalid JSON', () => {
emitMessage('INVALID'); emitMessage('INVALID');
jest.spyOn(socket, 'destroy'); jest.spyOn(socket, 'destroy');
+53 -26
View File
@@ -44,6 +44,8 @@ export class StratumV1Client {
private clientSuggestedDifficulty: SuggestDifficulty; private clientSuggestedDifficulty: SuggestDifficulty;
private stratumSubscription: Subscription; private stratumSubscription: Subscription;
private backgroundWork: NodeJS.Timeout[] = []; private backgroundWork: NodeJS.Timeout[] = [];
private readonly socketDataHandler: (data: Buffer) => void;
private destroyPromise: Promise<void> | null = null;
private statistics: StratumV1ClientStatistics; private statistics: StratumV1ClientStatistics;
private stratumInitialized = false; private stratumInitialized = false;
@@ -77,43 +79,68 @@ export class StratumV1Client {
private readonly redisMessagingService?: RedisMessagingService private readonly redisMessagingService?: RedisMessagingService
) { ) {
this.socket.on('data', (data: Buffer) => { this.socketDataHandler = (data: Buffer) => {
this.buffer += data.toString(); void this.handleSocketData(data);
let lines = this.buffer.split('\n'); };
this.buffer = lines.pop() || ''; // Save the last part of the data (incomplete line) to the buffer this.socket.on('data', this.socketDataHandler);
(async () => {
for (const m of lines.filter(l => l.length > 0)) {
if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) {
break;
}
try {
await this.handleMessage(m);
} catch (e) {
await this.socket.end();
console.error(e);
}
}
})();
});
} }
public async destroy() { public async destroy(): Promise<void> {
if (this.destroyPromise != null) {
if (this.clientEntity?.id) { return this.destroyPromise;
await this.redisMessagingService?.removeClientPresence(this.clientEntity.id, this.clientEntity.address);
await this.clientService.delete(this.clientEntity.id);
} }
this.destroyPromise = this.destroyInternal();
return this.destroyPromise;
}
private async destroyInternal(): Promise<void> {
this.connectionClosed = true;
this.socket.removeListener('data', this.socketDataHandler);
this.buffer = '';
if (this.stratumSubscription != null) { if (this.stratumSubscription != null) {
this.stratumSubscription.unsubscribe(); this.stratumSubscription.unsubscribe();
this.stratumSubscription = null;
} }
this.backgroundWork.forEach(work => { for (const work of this.backgroundWork) {
clearInterval(work); clearInterval(work);
}); }
this.backgroundWork = [];
this.miningSubmissionHashes.clear();
if (this.clientEntity?.id) {
const clientId = this.clientEntity.id;
const address = this.clientEntity.address;
this.clientEntity = null;
await this.redisMessagingService?.removeClientPresence(clientId, address);
await this.clientService.delete(clientId);
}
}
private async handleSocketData(data: Buffer): Promise<void> {
if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) {
return;
}
this.buffer += data.toString();
const lines = this.buffer.split('\n');
this.buffer = lines.pop() || ''; // Save the last part of the data (incomplete line) to the buffer
for (const m of lines.filter(l => l.length > 0)) {
if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) {
break;
}
try {
await this.handleMessage(m);
} catch (e) {
await this.socket.end();
console.error(e);
}
}
} }
private getRandomHexString() { private getRandomHexString() {
+54 -20
View File
@@ -125,20 +125,40 @@ export class StratumV1Service implements OnModuleInit {
let client: StratumV1Client | StratumV2Client = null; let client: StratumV1Client | StratumV2Client = null;
let protocol: 'v1' | 'v2' | null = null; let protocol: 'v1' | 'v2' | null = null;
let cleanedUp = false;
// Unified cleanup function // Unified cleanup function
const cleanup = async (reason: string) => { const cleanup = async (reason: string) => {
if (client != null && (protocol === 'v2' || (client as StratumV1Client).extraNonceAndSessionId != null)) { if (cleanedUp) {
await client.destroy(); return;
if (reason == 'Error') {
this.errorClosure++;
} else {
this.normalClosure++;
}
} }
if (!socket.destroyed) { cleanedUp = true;
socket.end();
socket.destroy(); const currentClient = client;
client = null;
try {
if (currentClient != null) {
const initializedClient = protocol === 'v2'
|| (currentClient as StratumV1Client).extraNonceAndSessionId != null;
await currentClient.destroy();
if (initializedClient) {
if (reason == 'Error') {
this.errorClosure++;
} else {
this.normalClosure++;
}
}
}
} finally {
socket.removeAllListeners('close');
socket.removeAllListeners('timeout');
socket.removeAllListeners('error');
socket.removeAllListeners('data');
if (!socket.destroyed) {
socket.end();
socket.destroy();
}
} }
}; };
@@ -241,19 +261,33 @@ export class StratumV1Service implements OnModuleInit {
socket.setKeepAlive(true, this.getTcpKeepAliveInitialDelayMs()); socket.setKeepAlive(true, this.getTcpKeepAliveInitialDelayMs());
const client = this.createV1Client(socket); const client = this.createV1Client(socket);
let cleanedUp = false;
const cleanup = async (reason: string) => { const cleanup = async (reason: string) => {
if (client.extraNonceAndSessionId != null) { if (cleanedUp) {
await client.destroy(); return;
if (reason === 'Error') {
this.errorClosure++;
} else {
this.normalClosure++;
}
} }
if (!socket.destroyed) { cleanedUp = true;
socket.end();
socket.destroy(); try {
const initializedClient = client.extraNonceAndSessionId != null;
await client.destroy();
if (initializedClient) {
if (reason === 'Error') {
this.errorClosure++;
} else {
this.normalClosure++;
}
}
} finally {
socket.removeAllListeners('close');
socket.removeAllListeners('timeout');
socket.removeAllListeners('error');
socket.removeAllListeners('data');
if (!socket.destroyed) {
socket.end();
socket.destroy();
}
} }
}; };