mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 17:15:03 -07:00
cleanup and support for all address types
This commit is contained in:
@@ -0,0 +1,83 @@
|
||||
import { Injectable } from '@nestjs/common';
|
||||
import { ConfigService } from '@nestjs/config';
|
||||
import { RPCClient } from 'rpc-bitcoin';
|
||||
import { BehaviorSubject, filter } from 'rxjs';
|
||||
|
||||
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
|
||||
import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo';
|
||||
|
||||
|
||||
|
||||
@Injectable()
|
||||
export class BitcoinRpcService {
|
||||
|
||||
private blockHeight = 0;
|
||||
private client: RPCClient;
|
||||
private _newBlock$: BehaviorSubject<IMiningInfo> = new BehaviorSubject(undefined);
|
||||
public newBlock$ = this._newBlock$.pipe(filter(block => block != null));
|
||||
|
||||
constructor(configService: ConfigService) {
|
||||
const url = configService.get('BITCOIN_RPC_URL');
|
||||
const user = configService.get('BITCOIN_RPC_USER');
|
||||
const pass = configService.get('BITCOIN_RPC_PASSWORD');
|
||||
const port = parseInt(configService.get('BITCOIN_RPC_PORT'));
|
||||
const timeout = parseInt(configService.get('BITCOIN_RPC_TIMEOUT'));
|
||||
|
||||
this.client = new RPCClient({ url, port, timeout, user, pass });
|
||||
|
||||
|
||||
console.log('Bitcoin RPC connected');
|
||||
|
||||
|
||||
// Maybe use ZeroMQ ?
|
||||
setInterval(async () => {
|
||||
const miningInfo = await this.getMiningInfo();
|
||||
if (miningInfo.blocks > this.blockHeight) {
|
||||
|
||||
this._newBlock$.next(miningInfo);
|
||||
|
||||
this.blockHeight = miningInfo.blocks;
|
||||
|
||||
}
|
||||
|
||||
}, 500);
|
||||
|
||||
}
|
||||
|
||||
|
||||
|
||||
|
||||
public async getBlockTemplate(): Promise<IBlockTemplate> {
|
||||
|
||||
const result: IBlockTemplate = await this.client.getblocktemplate({
|
||||
template_request: {
|
||||
rules: ['segwit'],
|
||||
mode: 'template',
|
||||
capabilities: ['serverlist', 'proposal']
|
||||
}
|
||||
});
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
public async getMiningInfo(): Promise<IMiningInfo> {
|
||||
return await this.client.getmininginfo();
|
||||
|
||||
}
|
||||
|
||||
public async SUBMIT_BLOCK(hexdata: string) {
|
||||
try {
|
||||
const res = await this.client.submitblock({
|
||||
hexdata
|
||||
});
|
||||
console.log(`BLOCK SUBMISSION RESPONSE: ${res}`);
|
||||
console.log(hexdata);
|
||||
console.log(JSON.stringify(res));
|
||||
process.exit();
|
||||
} catch (e) {
|
||||
console.log(`BLOCK SUBMISSION RESPONSE ERROR: ${e}`);
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
import { Injectable } from '@nestjs/common';
|
||||
import { from, map, Observable, shareReplay, switchMap, tap } from 'rxjs';
|
||||
|
||||
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
|
||||
import { BitcoinRpcService } from './bitcoin-rpc.service';
|
||||
|
||||
@Injectable()
|
||||
export class BlockTemplateService {
|
||||
|
||||
public currentBlockTemplate: IBlockTemplate;
|
||||
|
||||
public currentBlockTemplate$: Observable<{ blockTemplate: IBlockTemplate }>;
|
||||
|
||||
constructor(private readonly bitcoinRpcService: BitcoinRpcService) {
|
||||
this.currentBlockTemplate$ = this.bitcoinRpcService.newBlock$.pipe(
|
||||
switchMap((miningInfo) => from(this.bitcoinRpcService.getBlockTemplate()).pipe(map(blockTemplate => { return { miningInfo, blockTemplate } }))),
|
||||
tap(({ blockTemplate }) => this.currentBlockTemplate = blockTemplate),
|
||||
shareReplay({ refCount: true, bufferSize: 1 })
|
||||
);
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
import { Injectable, OnModuleInit } from '@nestjs/common';
|
||||
import { Cron, CronExpression } from '@nestjs/schedule';
|
||||
|
||||
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
|
||||
import { ClientService } from '../ORM/client/client.service';
|
||||
|
||||
@Injectable()
|
||||
export class CleanupService implements OnModuleInit {
|
||||
|
||||
constructor(
|
||||
private clientStatisticsService: ClientStatisticsService,
|
||||
private clientService: ClientService
|
||||
) {
|
||||
|
||||
}
|
||||
onModuleInit() {
|
||||
console.log('Cleanup service running.')
|
||||
}
|
||||
|
||||
@Cron(CronExpression.EVERY_HOUR)
|
||||
private async deleteOldStatistics() {
|
||||
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`);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
import { MiningJob } from '../models/MiningJob';
|
||||
|
||||
|
||||
export class StratumV1JobsService {
|
||||
|
||||
public latestJobId: number = 1;
|
||||
|
||||
public jobs: MiningJob[] = [];
|
||||
|
||||
|
||||
public addJob(job: MiningJob, clearJobs: boolean) {
|
||||
if (clearJobs) {
|
||||
this.jobs = [];
|
||||
}
|
||||
this.jobs.push(job);
|
||||
this.latestJobId++;
|
||||
}
|
||||
|
||||
public getLatestJob() {
|
||||
return this.jobs[this.jobs.length - 1];
|
||||
}
|
||||
|
||||
public getJobById(jobId: string) {
|
||||
return this.jobs.find(job => job.jobId == jobId);
|
||||
}
|
||||
|
||||
public getNextId() {
|
||||
return this.latestJobId.toString(16);
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
@@ -0,0 +1,74 @@
|
||||
import { Injectable, OnModuleInit } from '@nestjs/common';
|
||||
import { Server, Socket } from 'net';
|
||||
|
||||
import { StratumV1Client } from '../models/StratumV1Client';
|
||||
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
|
||||
import { ClientService } from '../ORM/client/client.service';
|
||||
import { BitcoinRpcService } from './bitcoin-rpc.service';
|
||||
import { BlockTemplateService } from './block-template.service';
|
||||
import { StratumV1JobsService } from './stratum-v1-jobs.service';
|
||||
|
||||
|
||||
@Injectable()
|
||||
export class StratumV1Service implements OnModuleInit {
|
||||
|
||||
// public clients: StratumV1Client[] = [];
|
||||
|
||||
constructor(
|
||||
private readonly bitcoinRpcService: BitcoinRpcService,
|
||||
private readonly blockTemplateService: BlockTemplateService,
|
||||
private readonly clientService: ClientService,
|
||||
private readonly clientStatisticsService: ClientStatisticsService
|
||||
) {
|
||||
}
|
||||
|
||||
|
||||
async onModuleInit(): Promise<void> {
|
||||
|
||||
//await this.clientStatisticsService.deleteAll();
|
||||
await this.clientService.deleteAll();
|
||||
|
||||
this.startSocketServer();
|
||||
|
||||
}
|
||||
|
||||
private startSocketServer() {
|
||||
new Server(async (socket: Socket) => {
|
||||
|
||||
|
||||
const client = new StratumV1Client(socket, new StratumV1JobsService(), this.blockTemplateService, this.bitcoinRpcService, this.clientService, this.clientStatisticsService);
|
||||
|
||||
|
||||
const clientCount = await this.clientService.connectedClientCount();
|
||||
|
||||
//this.clients.push(client);
|
||||
|
||||
console.log(`New client connected: ${socket.remoteAddress}, ${clientCount} total clients`);
|
||||
|
||||
socket.on('end', async () => {
|
||||
// Handle socket disconnection
|
||||
await this.clientService.delete(client.extraNonce);
|
||||
|
||||
const clientCount = await this.clientService.connectedClientCount();
|
||||
console.log(`Client disconnected: ${socket.remoteAddress}, ${clientCount} total clients`);
|
||||
});
|
||||
|
||||
socket.on('error', async (error: Error) => {
|
||||
|
||||
await this.clientService.delete(client.extraNonce);
|
||||
|
||||
const clientCount = await this.clientService.connectedClientCount();
|
||||
console.error(`Socket error:`, error);
|
||||
console.log(`Client disconnected: ${socket.remoteAddress}, ${clientCount} total clients`);
|
||||
|
||||
});
|
||||
|
||||
}).listen(3333, () => {
|
||||
console.log(`Bitcoin Stratum server is listening on port ${3333}`);
|
||||
});
|
||||
|
||||
}
|
||||
|
||||
|
||||
|
||||
}
|
||||
Reference in New Issue
Block a user