optimized share calculation and updating

This commit is contained in:
Ben Wilson
2024-02-17 18:22:09 -05:00
parent d486c778a3
commit e35bacc648
9 changed files with 161 additions and 84 deletions
+30 -16
View File
@@ -1,9 +1,9 @@
import { Injectable, OnModuleInit } from '@nestjs/common';
import { Interval } from '@nestjs/schedule';
import { DataSource } from 'typeorm';
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
import { ClientService } from '../ORM/client/client.service';
import { RpcBlockService } from '../ORM/rpc-block/rpc-block.service';
@Injectable()
export class AppService implements OnModuleInit {
@@ -11,7 +11,8 @@ export class AppService implements OnModuleInit {
constructor(
private readonly clientStatisticsService: ClientStatisticsService,
private readonly clientService: ClientService,
private readonly dataSource: DataSource
private readonly dataSource: DataSource,
private readonly rpcBlockService: RpcBlockService,
) {
}
@@ -28,25 +29,38 @@ export class AppService implements OnModuleInit {
await this.dataSource.query(`PRAGMA synchronous = off;`);
// //6Gb
// await this.dataSource.query(`PRAGMA mmap_size = 6000000000;`);
if (process.env.ENABLE_SOLO == 'true' && (process.env.NODE_APP_INSTANCE == null || process.env.NODE_APP_INSTANCE == '0')) {
setInterval(async () => {
await this.deleteOldStatistics();
}, 1000 * 60 * 60);
setInterval(async () => {
console.log('Killing dead clients');
await this.clientService.killDeadClients();
}, 1000 * 60 * 5);
setInterval(async () => {
console.log('Deleting Old Blocks');
await this.rpcBlockService.deleteOldBlocks();
}, 1000 * 60 * 60 * 24);
}
}
@Interval(1000 * 60 * 60)
private async deleteOldStatistics() {
console.log('Deleting statistics');
if (process.env.ENABLE_SOLO == 'true' && (process.env.NODE_APP_INSTANCE == null || process.env.NODE_APP_INSTANCE == '0')) {
const deletedStatistics = await this.clientStatisticsService.deleteOldStatistics();
console.log(`Deleted ${deletedStatistics.affected} old statistics`);
const deletedClients = await this.clientService.deleteOldClients();
console.log(`Deleted ${deletedClients.affected} old clients`);
}
const deletedStatistics = await this.clientStatisticsService.deleteOldStatistics();
console.log(`Deleted ${deletedStatistics.affected} old statistics`);
const deletedClients = await this.clientService.deleteOldClients();
console.log(`Deleted ${deletedClients.affected} old clients`);
}
@Interval(1000 * 60 * 5)
private async killDeadClients() {
console.log('Killing dead clients');
if (process.env.ENABLE_SOLO == 'true' && (process.env.NODE_APP_INSTANCE == null || process.env.NODE_APP_INSTANCE == '0')) {
await this.clientService.killDeadClients();
}
}
}
+29 -10
View File
@@ -1,15 +1,15 @@
import { Injectable } from '@nestjs/common';
import { Injectable, OnModuleInit } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { RPCClient } from 'rpc-bitcoin';
import { BehaviorSubject, filter, shareReplay } from 'rxjs';
import { RpcBlockService } from 'src/ORM/rpc-block/rpc-block.service';
import * as zmq from 'zeromq/v5-compat';
import * as zmq from 'zeromq';
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo';
@Injectable()
export class BitcoinRpcService {
export class BitcoinRpcService implements OnModuleInit {
private blockHeight = 0;
private client: RPCClient;
@@ -20,7 +20,9 @@ export class BitcoinRpcService {
private readonly configService: ConfigService,
private rpcBlockService: RpcBlockService
) {
}
async onModuleInit() {
const url = this.configService.get('BITCOIN_RPC_URL');
const user = this.configService.get('BITCOIN_RPC_USER');
const pass = this.configService.get('BITCOIN_RPC_PASSWORD');
@@ -36,19 +38,36 @@ export class BitcoinRpcService {
});
if (this.configService.get('BITCOIN_ZMQ_HOST')) {
const sock = zmq.socket("sub");
sock.connect(this.configService.get('BITCOIN_ZMQ_HOST'));
sock.subscribe("rawblock");
sock.on("message", async (topic: Buffer, message: Buffer) => {
console.log("new block zmq");
await this.pollMiningInfo();
console.log('Using ZMQ');
const sock = new zmq.Subscriber;
sock.connectTimeout = 1000;
sock.events.on('connect', () => {
console.log('ZMQ Connected');
});
this.pollMiningInfo().then(() => { });
sock.events.on('connect:retry', () => {
console.log('ZMQ Unable to connect, Retrying');
});
sock.connect(this.configService.get('BITCOIN_ZMQ_HOST'));
sock.subscribe('rawblock');
// Don't await this, otherwise it will block the rest of the program
this.listenForNewBlocks(sock);
await this.pollMiningInfo();
} else {
setInterval(this.pollMiningInfo.bind(this), 500);
}
}
private async listenForNewBlocks(sock: zmq.Subscriber) {
for await (const [topic, msg] of sock) {
console.log("New Block");
await this.pollMiningInfo();
}
}
public async pollMiningInfo() {
const miningInfo = await this.getMiningInfo();
if (miningInfo != null && miningInfo.blocks > this.blockHeight) {