data persistance init

This commit is contained in:
Ben Wilson
2023-06-23 22:30:54 -04:00
parent b963eaa514
commit 080471e52d
16 changed files with 2760 additions and 188 deletions
+3
View File
@@ -35,3 +35,6 @@ settings.json
!.vscode/tasks.json
!.vscode/launch.json
!.vscode/extensions.json
#DB
**.sqlite
+2461 -116
View File
File diff suppressed because it is too large Load Diff
+4 -1
View File
@@ -24,6 +24,7 @@
"@nestjs/config": "^2.3.2",
"@nestjs/core": "^9.0.0",
"@nestjs/platform-fastify": "^9.4.2",
"@nestjs/typeorm": "^10.0.0",
"big.js": "^6.2.1",
"bitcoinjs-lib": "^6.1.3",
"bs58": "^5.0.0",
@@ -32,7 +33,9 @@
"merkletreejs": "^0.3.10",
"reflect-metadata": "^0.1.13",
"rpc-bitcoin": "^2.0.0",
"rxjs": "^7.2.0"
"rxjs": "^7.2.0",
"sqlite3": "^5.1.6",
"typeorm": "^0.3.17"
},
"devDependencies": {
"@nestjs/cli": "^9.0.0",
@@ -0,0 +1,24 @@
import { Column, Entity, ManyToOne, PrimaryGeneratedColumn } from 'typeorm';
import { ClientEntity } from '../client/client.entity';
@Entity()
export class ClientStatisticsEntity {
@PrimaryGeneratedColumn()
id: number;
@Column({ type: 'datetime' })
time: Date;
@Column()
difficulty: number;
@ManyToOne(
() => ClientEntity,
client => client.clientStatistics
)
client: ClientEntity;
}
@@ -0,0 +1,14 @@
import { Global, Module } from '@nestjs/common';
import { TypeOrmModule } from '@nestjs/typeorm';
import { ClientStatisticsEntity } from './client-statistics.entity';
import { ClientStatisticsService } from './client-statistics.service';
@Global()
@Module({
imports: [TypeOrmModule.forFeature([ClientStatisticsEntity])],
providers: [ClientStatisticsService],
exports: [TypeOrmModule, ClientStatisticsService],
})
export class ClientStatisticsModule { }
@@ -0,0 +1,42 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { ClientStatisticsEntity } from './client-statistics.entity';
@Injectable()
export class ClientStatisticsService {
constructor(
@InjectRepository(ClientStatisticsEntity)
private clientStatisticsRepository: Repository<ClientStatisticsEntity>
) {
}
public async save(clientStatistic: Partial<ClientStatisticsEntity>) {
return await this.clientStatisticsRepository.save(clientStatistic);
}
public async getHashRate(clientId: string) {
const query = `
SELECT
(JULIANDAY(MAX(entry.time)) - JULIANDAY(MIN(entry.time))) * 24 * 60 * 60 AS timeDiff,
SUM(entry.difficulty) AS difficultySum
FROM
client_statistics_entity AS entry
WHERE
entry.clientId = ?
`;
const result = await this.clientStatisticsRepository.query(query, [clientId]);
const timeDiff = result[0].timeDiff;
const difficultySum = result[0].difficultySum;
return (difficultySum * 4294967296) / (timeDiff * 1000000000);
}
}
+29
View File
@@ -0,0 +1,29 @@
import { Column, Entity, OneToMany, PrimaryColumn } from 'typeorm';
import { ClientStatisticsEntity } from '../client-statistics/client-statistics.entity';
import { TrackedEntity } from '../utils/TrackedEntity.entity';
@Entity()
export class ClientEntity extends TrackedEntity {
@PrimaryColumn({ length: 8, type: 'varchar' })
id: string;
@Column({ type: 'datetime' })
startTime: Date;
@Column()
clientName: string;
@Column({ length: 62, type: 'varchar' })
address: string;
@Column({ default: 0 })
bestDifficulty: number
@OneToMany(
() => ClientStatisticsEntity,
clientStatistic => clientStatistic.client
)
clientStatistics?: ClientStatisticsEntity[];
}
+14
View File
@@ -0,0 +1,14 @@
import { Global, Module } from '@nestjs/common';
import { TypeOrmModule } from '@nestjs/typeorm';
import { ClientEntity } from './client.entity';
import { ClientService } from './client.service';
@Global()
@Module({
imports: [TypeOrmModule.forFeature([ClientEntity])],
providers: [ClientService],
exports: [TypeOrmModule, ClientService],
})
export class ClientModule { }
+51
View File
@@ -0,0 +1,51 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { ClientEntity } from './client.entity';
@Injectable()
export class ClientService {
constructor(
@InjectRepository(ClientEntity)
private clientRepository: Repository<ClientEntity>
) {
}
public async save(client: Partial<ClientEntity>) {
return await this.clientRepository.save(client);
}
public async delete(id: string) {
return await this.clientRepository.softDelete({ id });
}
public async updateBestDifficulty(id: string, bestDifficulty: number) {
return await this.clientRepository.update(id, { bestDifficulty });
}
public async connectedClientCount(): Promise<number> {
return await this.clientRepository.count();
}
public async getByAddress(address: string): Promise<ClientEntity[]> {
return await this.clientRepository.find({
where: {
address
}
})
}
public async getByAddressAndName(address: string, id: string): Promise<ClientEntity> {
return await this.clientRepository.findOne({
where: {
address,
id
}
})
}
}
+12
View File
@@ -0,0 +1,12 @@
import { CreateDateColumn, DeleteDateColumn, UpdateDateColumn } from 'typeorm';
export abstract class TrackedEntity {
@DeleteDateColumn({ nullable: true, type: 'datetime' })
public deletedAt?: Date;
@CreateDateColumn({ type: 'datetime' })
public createdAt?: Date
@UpdateDateColumn({ type: 'datetime' })
public updatedAt?: Date
}
+34 -19
View File
@@ -1,44 +1,59 @@
import { Controller, Get, Param } from '@nestjs/common';
import { AppService } from './app.service';
import { ClientStatisticsService } from './ORM/client-statistics/client-statistics.service';
import { ClientService } from './ORM/client/client.service';
import { StratumV1Service } from './stratum-v1.service';
@Controller()
export class AppController {
constructor(private readonly appService: AppService, private readonly stratumV1Service: StratumV1Service) { }
constructor(
private readonly appService: AppService,
private readonly stratumV1Service: StratumV1Service,
private readonly clientService: ClientService,
private readonly clientStatisticsService: ClientStatisticsService
) { }
@Get()
getInfo() {
async getInfo() {
return {
clients: this.stratumV1Service.clients.length
clients: await this.clientService.connectedClientCount()
}
}
@Get('client/:clientId')
getClientInfo(@Param('clientId') clientId: string) {
const workers = this.stratumV1Service.clients.filter(client => client.clientAuthorization.address === clientId);
async getClientInfo(@Param('clientId') clientId: string) {
const workers = await this.clientService.getByAddress(clientId);
//const workers = this.stratumV1Service.clients.filter(client => client.clientAuthorization.address === clientId);
return {
workersCount: workers.length,
workers: workers.map(worker => {
return {
id: worker.id,
name: worker.clientAuthorization.worker,
bestDifficulty: Math.floor(worker.statistics.bestDifficulty),
hashRate: Math.floor(worker.statistics.getHashRate()),
startTime: worker.startTime
}
})
workers: await Promise.all(
workers.map(async (worker) => {
return {
id: worker.id,
name: worker.clientName,
bestDifficulty: Math.floor(worker.bestDifficulty),
hashRate: Math.floor(await this.clientStatisticsService.getHashRate(worker.id)),
startTime: worker.startTime
};
})
)
}
}
@Get('client/:clientId/:workerId')
getWorkerInfo(@Param('clientId') clientId: string, @Param('workerId') workerId: string) {
const worker = this.stratumV1Service.clients.find(client => client.clientAuthorization.address === clientId && client.id === workerId);
async getWorkerInfo(@Param('clientId') clientId: string, @Param('workerId') workerId: string) {
const worker = await this.clientService.getByAddressAndName(clientId, workerId);
//const worker = this.stratumV1Service.clients.find(client => client.clientAuthorization.address === clientId && client.id === workerId);
return {
id: worker.id,
name: worker.clientAuthorization.worker,
bestDifficulty: Math.floor(worker.statistics.bestDifficulty),
hashData: worker.statistics.historicSubmissions,
name: worker.clientName,
bestDifficulty: Math.floor(worker.bestDifficulty),
//hashData: worker.statistics.historicSubmissions,
startTime: worker.startTime
}
}
+15 -2
View File
@@ -1,17 +1,30 @@
import { Module } from '@nestjs/common';
import { ConfigModule } from '@nestjs/config';
import { TypeOrmModule } from '@nestjs/typeorm';
import { AppController } from './app.controller';
import { AppService } from './app.service';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { BlockTemplateService } from './BlockTemplateService';
import { ClientStatisticsModule } from './ORM/client-statistics/client-statistics.module';
import { ClientModule } from './ORM/client/client.module';
import { StratumV1Service } from './stratum-v1.service';
const ORMModules = [
ClientStatisticsModule,
ClientModule
]
@Module({
imports: [
ConfigModule.forRoot()
ConfigModule.forRoot(),
TypeOrmModule.forRoot({
type: 'sqlite',
database: './DB/public-pool.sqlite',
synchronize: true,
autoLoadEntities: true,
}),
...ORMModules
],
controllers: [AppController],
providers: [
+2 -2
View File
@@ -70,8 +70,8 @@ export class BitcoinRpcService {
hexdata
});
console.log(`BLOCK SUBMISSION RESPONSE: ${res}`);
// console.log(JSON.stringify(res));
// process.exit();
console.log(JSON.stringify(res));
process.exit();
} catch (e) {
console.log(`BLOCK SUBMISSION RESPONSE ERROR: ${e}`);
}
+28 -7
View File
@@ -8,6 +8,9 @@ import { combineLatest, interval, startWith } from 'rxjs';
import { BitcoinRpcService } from '../bitcoin-rpc.service';
import { BlockTemplateService } from '../BlockTemplateService';
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
import { ClientEntity } from '../ORM/client/client.entity';
import { ClientService } from '../ORM/client/client.service';
import { StratumV1JobsService } from '../stratum-v1-jobs.service';
import { EasyUnsubscribe } from '../utils/AutoUnsubscribe';
import { eRequestMethod } from './enums/eRequestMethod';
@@ -22,12 +25,14 @@ import { StratumV1ClientStatistics } from './StratumV1ClientStatistics';
export class StratumV1Client extends EasyUnsubscribe {
public startTime: Date;
public clientSubscription: SubscriptionMessage;
public clientConfiguration: ConfigurationMessage;
public clientAuthorization: AuthorizationMessage;
public clientSuggestedDifficulty: SuggestDifficulty;
public statistics: StratumV1ClientStatistics = new StratumV1ClientStatistics();
public statistics: StratumV1ClientStatistics;
public id: string;
public stratumInitialized = false;
@@ -38,16 +43,20 @@ export class StratumV1Client extends EasyUnsubscribe {
public jobRefreshInterval: NodeJS.Timer;
public entity: ClientEntity;
constructor(
public readonly socket: Socket,
private readonly stratumV1JobsService: StratumV1JobsService,
private readonly blockTemplateService: BlockTemplateService,
private readonly bitcoinRpcService: BitcoinRpcService
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly clientService: ClientService,
private readonly clientStatisticsService: ClientStatisticsService
) {
super();
this.startTime = new Date();
this.statistics = new StratumV1ClientStatistics(this.clientStatisticsService);
this.id = this.getRandomHexString();
console.log(`New client ID: : ${this.id}`);
@@ -194,7 +203,7 @@ export class StratumV1Client extends EasyUnsubscribe {
const errors = await validate(miningSubmitMessage, validatorOptions);
if (errors.length === 0) {
this.handleMiningSubmission(miningSubmitMessage);
await this.handleMiningSubmission(miningSubmitMessage);
socket.write(JSON.stringify(miningSubmitMessage.response()) + '\n');
} else {
console.error(errors);
@@ -213,6 +222,15 @@ export class StratumV1Client extends EasyUnsubscribe {
this.stratumInitialized = true;
this.entity = await this.clientService.save({
id: this.id,
address: this.clientAuthorization.address,
clientName: this.clientAuthorization.worker,
startTime: new Date(),
});
let lastIntervalCount = undefined;
combineLatest([this.blockTemplateService.currentBlockTemplate$, interval(60000).pipe(startWith(-1))]).subscribe(([{ miningInfo, blockTemplate }, interValCount]) => {
let clearJobs = false;
@@ -235,7 +253,7 @@ export class StratumV1Client extends EasyUnsubscribe {
}
private handleMiningSubmission(submission: MiningSubmitMessage) {
private async handleMiningSubmission(submission: MiningSubmitMessage) {
const job = this.stratumV1JobsService.getJobById(submission.jobId);
// a miner may submit a job that doesn't exist anymore if it was removed by a new block notification
@@ -247,8 +265,11 @@ export class StratumV1Client extends EasyUnsubscribe {
if (diff >= this.clientDifficulty) {
const networkDifficulty = this.calculateNetworkDifficulty(parseInt(job.blockTemplate.bits, 16));
this.statistics.addSubmission(this.clientDifficulty, diff, networkDifficulty);
if (diff >= networkDifficulty) {
await this.statistics.addSubmission(this.entity, this.clientDifficulty);
if (diff > this.entity.bestDifficulty) {
await this.clientService.updateBestDifficulty(this.id, diff);
}
if (diff >= (networkDifficulty / 2)) {
console.log('!!! BOCK FOUND !!!');
this.constructBlockAndBroadcast(job, submission);
}
+8 -31
View File
@@ -1,45 +1,22 @@
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
import { ClientEntity } from '../ORM/client/client.entity';
export class StratumV1ClientStatistics {
public bestDifficulty: number = 0;
public historicSubmissions = [];
constructor() {
constructor(private readonly clientStatisticsService: ClientStatisticsService) {
}
public addSubmission(targetDifficulty: number, submissionDifficulty: number, networkDifficulty: number) {
public async addSubmission(client: ClientEntity, targetDifficulty: number) {
if (this.historicSubmissions.length >= 1000) {
this.historicSubmissions.shift();
}
this.historicSubmissions.push({
await this.clientStatisticsService.save({
time: new Date(),
difficulty: targetDifficulty,
time: new Date()
client
});
if (submissionDifficulty > this.bestDifficulty) {
this.bestDifficulty = submissionDifficulty;
}
}
public getHashRate() {
return this.historicSubmissions.reduce((pre, cur, idx, arr) => {
if (idx === 0) {
pre.time = cur.time;
}
pre.difficulty += cur.difficulty;
if (arr.length - 1 === idx) {
const duration = cur.time.getTime() - pre.time.getTime();
const sumDifficulty = pre.difficulty;
return (sumDifficulty * 4294967296) / (duration * 1000000)
}
return pre;
}, {
difficulty: 0,
time: undefined
})
}
}
+18 -9
View File
@@ -4,17 +4,21 @@ import { Server, Socket } from 'net';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { BlockTemplateService } from './BlockTemplateService';
import { StratumV1Client } from './models/StratumV1Client';
import { ClientStatisticsService } from './ORM/client-statistics/client-statistics.service';
import { ClientService } from './ORM/client/client.service';
import { StratumV1JobsService } from './stratum-v1-jobs.service';
@Injectable()
export class StratumV1Service implements OnModuleInit {
public clients: StratumV1Client[] = [];
// public clients: StratumV1Client[] = [];
constructor(
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly blockTemplateService: BlockTemplateService
private readonly blockTemplateService: BlockTemplateService,
private readonly clientService: ClientService,
private readonly clientStatisticsService: ClientStatisticsService
) {
}
@@ -26,19 +30,24 @@ export class StratumV1Service implements OnModuleInit {
}
private startSocketServer() {
new Server((socket: Socket) => {
new Server(async (socket: Socket) => {
const client = new StratumV1Client(socket, new StratumV1JobsService(), this.blockTemplateService, this.bitcoinRpcService, this.clientService, this.clientStatisticsService);
const client = new StratumV1Client(socket, new StratumV1JobsService(), this.blockTemplateService, this.bitcoinRpcService);
this.clients.push(client);
console.log(`New client connected: ${socket.remoteAddress}, ${this.clients.length} total clients`);
const clientCount = await this.clientService.connectedClientCount();
socket.on('end', () => {
//this.clients.push(client);
console.log(`New client connected: ${socket.remoteAddress}, ${clientCount} total clients`);
socket.on('end', async () => {
// Handle socket disconnection
this.clients = this.clients.filter(c => c.id == client.id);
console.log(`Client disconnected: ${socket.remoteAddress}, ${this.clients.length} total clients`);
await this.clientService.delete(client.id);
const clientCount = await this.clientService.connectedClientCount();
console.log(`Client disconnected: ${socket.remoteAddress}, ${clientCount} total clients`);
});
socket.on('error', (error: Error) => {