From 6f6abe5e2db2bafc2ff22ef9ae45334851a8d16d Mon Sep 17 00:00:00 2001 From: Ben Date: Fri, 21 Nov 2025 18:42:00 -0500 Subject: [PATCH] batch kill dead clients --- src/ORM/client/client.service.ts | 38 +++++++++++++++++++++++--------- src/services/app.service.ts | 4 +++- 2 files changed, 31 insertions(+), 11 deletions(-) diff --git a/src/ORM/client/client.service.ts b/src/ORM/client/client.service.ts index ec29925..4628c27 100644 --- a/src/ORM/client/client.service.ts +++ b/src/ORM/client/client.service.ts @@ -9,7 +9,7 @@ import { ClientEntity } from './client.entity'; @Injectable() export class ClientService { - private heartbeatBulkUpdate: { [id: string]:{id: string, hashRate: number, updatedAt: Date}} = {}; + private heartbeatBulkUpdate: { [id: string]: { id: string, hashRate: number, updatedAt: Date } } = {}; constructor( @InjectRepository(ClientEntity) @@ -19,36 +19,54 @@ export class ClientService { } public async killDeadClients() { + const batchSize = 500; - return await this.clientRepository + const clients = await this.clientRepository + .createQueryBuilder('client') + .where('client.deletedAt IS NULL') + .andWhere('client.updatedAt < NOW() - INTERVAL \'5 minutes\'') + .orderBy('client.id', 'ASC') // consistent lock order ← prevents deadlocks + .limit(batchSize) // small batches ← prevents lock explosion + .setLock('pessimistic_write') // adds FOR UPDATE SKIP LOCKED + .getMany(); + + if (clients.length === 0) { + return false; // nothing to do (caller can stop looping) + } + + const ids = clients.map(c => c.id); + + await this.clientRepository .createQueryBuilder() .update(ClientEntity) - .set({ deletedAt: () => "NOW()" }) - .where("deletedAt IS NULL AND updatedAt < NOW() - interval '5 minutes' ") + .set({ deletedAt: () => 'NOW()' }) + .whereInIds(ids) .execute(); + + return true; // tell caller there might be more batches } //public async heartbeat(id, hashRate: number, updatedAt: Date) { // return await this.clientRepository.update({ id }, { hashRate, deletedAt: null, updatedAt }); // } - public heartbeatBulkAsync(id, hashRate: number, updatedAt: Date){ - if(this.heartbeatBulkUpdate[id] != null){ + public heartbeatBulkAsync(id, hashRate: number, updatedAt: Date) { + if (this.heartbeatBulkUpdate[id] != null) { this.heartbeatBulkUpdate[id].hashRate = hashRate; this.heartbeatBulkUpdate[id].updatedAt = updatedAt; return; } - this.heartbeatBulkUpdate[id] = {id, hashRate, updatedAt}; + this.heartbeatBulkUpdate[id] = { id, hashRate, updatedAt }; } - public async doBulkHeartbeatUpdate(){ - if(Object.keys(this.heartbeatBulkUpdate).length < 1){ + public async doBulkHeartbeatUpdate() { + if (Object.keys(this.heartbeatBulkUpdate).length < 1) { console.log('No heartbeats to update.') return; } const values = Object.entries(this.heartbeatBulkUpdate).map(([key, value]) => { - return `('${value.id}', ${value.hashRate}, NOW())` + return `('${value.id}', ${value.hashRate}, NOW())` }).join(','); const query = ` diff --git a/src/services/app.service.ts b/src/services/app.service.ts index 1c8fd8f..238d4d7 100644 --- a/src/services/app.service.ts +++ b/src/services/app.service.ts @@ -30,7 +30,9 @@ export class AppService implements OnModuleInit { setInterval(async () => { console.log('Killing dead clients'); - await this.clientService.killDeadClients(); + while (await this.clientService.killDeadClients()) { + console.log('Finished killing clients'); + } }, 1000 * 60 * 5); setInterval(async () => {