Compare commits

3 Commits
Author SHA1 Message Date
Ben efbd95e464 Keep SV1 client entity during cleanup 2026-06-09 12:45:33 -04:00
Ben bb0acf0395 Pass Stratum backpressure setting to container 2026-06-09 12:42:27 -04:00
Ben 9b99c5153a Harden Stratum V1 cleanup 2026-06-09 12:30:03 -04:00
6 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
@@ -69,6 +69,7 @@ services:
SECURE_STRATUM_PORTS: ${SECURE_STRATUM_PORTS:-4333,4332,4331,4330} SECURE_STRATUM_PORTS: ${SECURE_STRATUM_PORTS:-4333,4332,4331,4330}
STRATUM_MAX_CONNECTIONS_PER_LISTENER: ${STRATUM_MAX_CONNECTIONS_PER_LISTENER:-10000} STRATUM_MAX_CONNECTIONS_PER_LISTENER: ${STRATUM_MAX_CONNECTIONS_PER_LISTENER:-10000}
STRATUM_TLS_HANDSHAKE_TIMEOUT_MS: ${STRATUM_TLS_HANDSHAKE_TIMEOUT_MS:-10000} STRATUM_TLS_HANDSHAKE_TIMEOUT_MS: ${STRATUM_TLS_HANDSHAKE_TIMEOUT_MS:-10000}
STRATUM_BACKPRESSURE_ENABLED: "${STRATUM_BACKPRESSURE_ENABLED:-false}"
STRATUM_V2_PORTS: ${STRATUM_V2_PORTS:-} STRATUM_V2_PORTS: ${STRATUM_V2_PORTS:-}
SHARE_ACCOUNTING_BATCH_SIZE: ${SHARE_ACCOUNTING_BATCH_SIZE:-500} SHARE_ACCOUNTING_BATCH_SIZE: ${SHARE_ACCOUNTING_BATCH_SIZE:-500}
SHARE_ACCOUNTING_FLUSH_INTERVAL_MS: ${SHARE_ACCOUNTING_FLUSH_INTERVAL_MS:-25} SHARE_ACCOUNTING_FLUSH_INTERVAL_MS: ${SHARE_ACCOUNTING_FLUSH_INTERVAL_MS:-25}
+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');
+49 -23
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,12 +79,56 @@ export class StratumV1Client {
private readonly redisMessagingService?: RedisMessagingService private readonly redisMessagingService?: RedisMessagingService
) { ) {
this.socket.on('data', (data: Buffer) => { this.socketDataHandler = (data: Buffer) => {
void this.handleSocketData(data);
};
this.socket.on('data', this.socketDataHandler);
}
public async destroy(): Promise<void> {
if (this.destroyPromise != null) {
return this.destroyPromise;
}
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) {
this.stratumSubscription.unsubscribe();
this.stratumSubscription = null;
}
for (const work of this.backgroundWork) {
clearInterval(work);
}
this.backgroundWork = [];
this.miningSubmissionHashes.clear();
if (this.clientEntity?.id) {
const clientId = this.clientEntity.id;
const address = this.clientEntity.address;
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(); this.buffer += data.toString();
let lines = this.buffer.split('\n'); const lines = this.buffer.split('\n');
this.buffer = lines.pop() || ''; // Save the last part of the data (incomplete line) to the buffer this.buffer = lines.pop() || ''; // Save the last part of the data (incomplete line) to the buffer
(async () => {
for (const m of lines.filter(l => l.length > 0)) { for (const m of lines.filter(l => l.length > 0)) {
if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) { if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) {
break; break;
@@ -94,26 +140,6 @@ export class StratumV1Client {
console.error(e); console.error(e);
} }
} }
})();
});
}
public async destroy() {
if (this.clientEntity?.id) {
await this.redisMessagingService?.removeClientPresence(this.clientEntity.id, this.clientEntity.address);
await this.clientService.delete(this.clientEntity.id);
}
if (this.stratumSubscription != null) {
this.stratumSubscription.unsubscribe();
}
this.backgroundWork.forEach(work => {
clearInterval(work);
});
} }
private getRandomHexString() { private getRandomHexString() {
+37 -3
View File
@@ -125,21 +125,41 @@ 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;
}
cleanedUp = true;
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') { if (reason == 'Error') {
this.errorClosure++; this.errorClosure++;
} else { } else {
this.normalClosure++; this.normalClosure++;
} }
} }
}
} finally {
socket.removeAllListeners('close');
socket.removeAllListeners('timeout');
socket.removeAllListeners('error');
socket.removeAllListeners('data');
if (!socket.destroyed) { if (!socket.destroyed) {
socket.end(); socket.end();
socket.destroy(); socket.destroy();
} }
}
}; };
// Handle client disconnection // Handle client disconnection
@@ -241,20 +261,34 @@ 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) {
return;
}
cleanedUp = true;
try {
const initializedClient = client.extraNonceAndSessionId != null;
await client.destroy(); await client.destroy();
if (initializedClient) {
if (reason === 'Error') { if (reason === 'Error') {
this.errorClosure++; this.errorClosure++;
} else { } else {
this.normalClosure++; this.normalClosure++;
} }
} }
} finally {
socket.removeAllListeners('close');
socket.removeAllListeners('timeout');
socket.removeAllListeners('error');
socket.removeAllListeners('data');
if (!socket.destroyed) { if (!socket.destroyed) {
socket.end(); socket.end();
socket.destroy(); socket.destroy();
} }
}
}; };
socket.on('close', async (hadError: boolean) => { socket.on('close', async (hadError: boolean) => {