mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 09:05:06 -07:00
killDeadClients
This commit is contained in:
@@ -6,6 +6,10 @@ import { TrackedEntity } from '../utils/TrackedEntity.entity';
|
|||||||
|
|
||||||
@Entity()
|
@Entity()
|
||||||
@Index("IDX_unique_nonce", { synchronize: false })
|
@Index("IDX_unique_nonce", { synchronize: false })
|
||||||
|
@Index('idx_client_cleanup', ['updatedAt', 'id'], {
|
||||||
|
where: '"deletedAt" IS NULL',
|
||||||
|
// This is a partial (filtered) index — only indexes active clients
|
||||||
|
})
|
||||||
export class ClientEntity extends TrackedEntity {
|
export class ClientEntity extends TrackedEntity {
|
||||||
|
|
||||||
@PrimaryGeneratedColumn('uuid')
|
@PrimaryGeneratedColumn('uuid')
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
import { Injectable } from '@nestjs/common';
|
import { Injectable } from '@nestjs/common';
|
||||||
import { InjectRepository } from '@nestjs/typeorm';
|
import { InjectDataSource, InjectRepository } from '@nestjs/typeorm';
|
||||||
import { Repository } from 'typeorm';
|
import { DataSource, Repository } from 'typeorm';
|
||||||
|
|
||||||
import { ClientEntity } from './client.entity';
|
import { ClientEntity } from './client.entity';
|
||||||
|
|
||||||
@@ -12,38 +12,37 @@ 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(
|
||||||
|
@InjectDataSource()
|
||||||
|
private dataSource: DataSource,
|
||||||
@InjectRepository(ClientEntity)
|
@InjectRepository(ClientEntity)
|
||||||
private clientRepository: Repository<ClientEntity>
|
private clientRepository: Repository<ClientEntity>
|
||||||
) {
|
) {
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
public async killDeadClients() {
|
public async killDeadClients(): Promise<boolean> {
|
||||||
const batchSize = 500;
|
return this.dataSource.transaction(async (manager) => {
|
||||||
|
const clients = await manager
|
||||||
|
.getRepository(ClientEntity)
|
||||||
|
.createQueryBuilder('c')
|
||||||
|
.where('c.deletedAt IS NULL')
|
||||||
|
.andWhere('c.updatedAt < NOW() - INTERVAL \'5 minutes\'')
|
||||||
|
.orderBy('c.id')
|
||||||
|
.limit(500)
|
||||||
|
.setLock('pessimistic_write') // FOR UPDATE SKIP LOCKED
|
||||||
|
.getMany();
|
||||||
|
|
||||||
const clients = await this.clientRepository
|
if (clients.length === 0) return false;
|
||||||
.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) {
|
await manager
|
||||||
return false; // nothing to do (caller can stop looping)
|
.createQueryBuilder()
|
||||||
}
|
.update(ClientEntity)
|
||||||
|
.set({ deletedAt: () => 'NOW()' })
|
||||||
|
.whereInIds(clients.map(c => c.id))
|
||||||
|
.execute();
|
||||||
|
|
||||||
const ids = clients.map(c => c.id);
|
return true;
|
||||||
|
});
|
||||||
await this.clientRepository
|
|
||||||
.createQueryBuilder()
|
|
||||||
.update(ClientEntity)
|
|
||||||
.set({ deletedAt: () => 'NOW()' })
|
|
||||||
.whereInIds(ids)
|
|
||||||
.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) {
|
||||||
|
|||||||
Reference in New Issue
Block a user