mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 17:15:03 -07:00
batch kill dead clients
This commit is contained in:
@@ -9,7 +9,7 @@ import { ClientEntity } from './client.entity';
|
|||||||
@Injectable()
|
@Injectable()
|
||||||
export class ClientService {
|
export class ClientService {
|
||||||
|
|
||||||
private heartbeatBulkUpdate: { [id: string]:{id: string, hashRate: number, updatedAt: Date}} = {};
|
private heartbeatBulkUpdate: { [id: string]: { id: string, hashRate: number, updatedAt: Date } } = {};
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
@InjectRepository(ClientEntity)
|
@InjectRepository(ClientEntity)
|
||||||
@@ -19,30 +19,48 @@ export class ClientService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
public async killDeadClients() {
|
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()
|
.createQueryBuilder()
|
||||||
.update(ClientEntity)
|
.update(ClientEntity)
|
||||||
.set({ deletedAt: () => "NOW()" })
|
.set({ deletedAt: () => 'NOW()' })
|
||||||
.where("deletedAt IS NULL AND updatedAt < NOW() - interval '5 minutes' ")
|
.whereInIds(ids)
|
||||||
.execute();
|
.execute();
|
||||||
|
|
||||||
|
return true; // tell caller there might be more batches
|
||||||
}
|
}
|
||||||
|
|
||||||
//public async heartbeat(id, hashRate: number, updatedAt: Date) {
|
//public async heartbeat(id, hashRate: number, updatedAt: Date) {
|
||||||
// return await this.clientRepository.update({ id }, { hashRate, deletedAt: null, updatedAt });
|
// return await this.clientRepository.update({ id }, { hashRate, deletedAt: null, updatedAt });
|
||||||
// }
|
// }
|
||||||
|
|
||||||
public heartbeatBulkAsync(id, hashRate: number, updatedAt: Date){
|
public heartbeatBulkAsync(id, hashRate: number, updatedAt: Date) {
|
||||||
if(this.heartbeatBulkUpdate[id] != null){
|
if (this.heartbeatBulkUpdate[id] != null) {
|
||||||
this.heartbeatBulkUpdate[id].hashRate = hashRate;
|
this.heartbeatBulkUpdate[id].hashRate = hashRate;
|
||||||
this.heartbeatBulkUpdate[id].updatedAt = updatedAt;
|
this.heartbeatBulkUpdate[id].updatedAt = updatedAt;
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
this.heartbeatBulkUpdate[id] = {id, hashRate, updatedAt};
|
this.heartbeatBulkUpdate[id] = { id, hashRate, updatedAt };
|
||||||
}
|
}
|
||||||
|
|
||||||
public async doBulkHeartbeatUpdate(){
|
public async doBulkHeartbeatUpdate() {
|
||||||
if(Object.keys(this.heartbeatBulkUpdate).length < 1){
|
if (Object.keys(this.heartbeatBulkUpdate).length < 1) {
|
||||||
console.log('No heartbeats to update.')
|
console.log('No heartbeats to update.')
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -30,7 +30,9 @@ export class AppService implements OnModuleInit {
|
|||||||
|
|
||||||
setInterval(async () => {
|
setInterval(async () => {
|
||||||
console.log('Killing dead clients');
|
console.log('Killing dead clients');
|
||||||
await this.clientService.killDeadClients();
|
while (await this.clientService.killDeadClients()) {
|
||||||
|
console.log('Finished killing clients');
|
||||||
|
}
|
||||||
}, 1000 * 60 * 5);
|
}, 1000 * 60 * 5);
|
||||||
|
|
||||||
setInterval(async () => {
|
setInterval(async () => {
|
||||||
|
|||||||
Reference in New Issue
Block a user