diff --git a/ecosystem.config.js b/ecosystem.config.js index 69db226..f14b923 100644 --- a/ecosystem.config.js +++ b/ecosystem.config.js @@ -8,16 +8,18 @@ module.exports = { env: { MASTER: 'true', }, + time: true }, // Worker instances { name: 'workers', script: './dist/main.js', - instances: 31, + instances: 2, exec_mode: "cluster", env: { MASTER: 'false', }, + time: true }, ], }; \ No newline at end of file diff --git a/package-lock.json b/package-lock.json index 53b032b..12c07b3 100644 --- a/package-lock.json +++ b/package-lock.json @@ -30,6 +30,7 @@ "merkle-lib": "^2.0.10", "node-telegram-bot-api": "^0.61.0", "pg": "^8.11.3", + "pg-pubsub": "^0.8.1", "reflect-metadata": "^0.1.13", "rpc-bitcoin": "^2.0.0", "rxjs": "^7.2.0", @@ -8852,6 +8853,14 @@ "resolved": "https://registry.npmjs.org/pg-connection-string/-/pg-connection-string-2.6.2.tgz", "integrity": "sha512-ch6OwaeaPYcova4kKZ15sbJ2hKb/VP48ZD2gE7i1J+L4MspCtBMAx8nMgz7bksc7IojCIIWuEhHibSMFH8m8oA==" }, + "node_modules/pg-format": { + "version": "1.0.4", + "resolved": "https://registry.npmjs.org/pg-format/-/pg-format-1.0.4.tgz", + "integrity": "sha512-YyKEF78pEA6wwTAqOUaHIN/rWpfzzIuMh9KdAhc3rSLQ/7zkRFcCgYBAEGatDstLyZw4g0s9SNICmaTGnBVeyw==", + "engines": { + "node": ">=4.0" + } + }, "node_modules/pg-int8": { "version": "1.0.1", "resolved": "https://registry.npmjs.org/pg-int8/-/pg-int8-1.0.1.tgz", @@ -8873,6 +8882,20 @@ "resolved": "https://registry.npmjs.org/pg-protocol/-/pg-protocol-1.6.0.tgz", "integrity": "sha512-M+PDm637OY5WM307051+bsDia5Xej6d9IR4GwJse1qA1DIhiKlksvrneZOYQq42OM+spubpcNYEo2FcKQrDk+Q==" }, + "node_modules/pg-pubsub": { + "version": "0.8.1", + "resolved": "https://registry.npmjs.org/pg-pubsub/-/pg-pubsub-0.8.1.tgz", + "integrity": "sha512-b/EHOwCrag4isghc4XgRipeAjfgyNg1DiL3Dwwh1Ojp91Lriltn5eg2nSWjBe4pzcFzhTM6HiB7LOG9NN1nx5g==", + "dependencies": { + "pg": "^8.7.3", + "pg-format": "^1.0.2", + "pony-cause": "^2.1.8", + "promised-retry": "^0.5.0" + }, + "engines": { + "node": "^14.18.0 || >=16.0.0" + } + }, "node_modules/pg-types": { "version": "2.2.0", "resolved": "https://registry.npmjs.org/pg-types/-/pg-types-2.2.0.tgz", @@ -9031,6 +9054,14 @@ "node": ">=4" } }, + "node_modules/pony-cause": { + "version": "2.1.11", + "resolved": "https://registry.npmjs.org/pony-cause/-/pony-cause-2.1.11.tgz", + "integrity": "sha512-M7LhCsdNbNgiLYiP4WjsfLUuFmCfnjdF6jKe2R9NKl4WFN+HZPGHJZ9lnLP7f9ZnKe3U9nuWD0szirmj+migUg==", + "engines": { + "node": ">=12.0.0" + } + }, "node_modules/postgres-array": { "version": "2.0.0", "resolved": "https://registry.npmjs.org/postgres-array/-/postgres-array-2.0.0.tgz", @@ -9165,6 +9196,25 @@ "node": ">=10" } }, + "node_modules/promised-retry": { + "version": "0.5.0", + "resolved": "https://registry.npmjs.org/promised-retry/-/promised-retry-0.5.0.tgz", + "integrity": "sha512-jbYvN6UGE+/3E1g0JmgDPchUc+4VI4cBaPjdr2Lso22xfFqut2warEf6IhWuhPJKbJYVOQAyCt2Jx+01ORCItg==", + "dependencies": { + "pony-cause": "^1.1.1" + }, + "engines": { + "node": "^14.17.0 || >=16.0.0" + } + }, + "node_modules/promised-retry/node_modules/pony-cause": { + "version": "1.1.1", + "resolved": "https://registry.npmjs.org/pony-cause/-/pony-cause-1.1.1.tgz", + "integrity": "sha512-PxkIc/2ZpLiEzQXu5YRDOUgBlfGYBY8156HY5ZcRAwwonMk5W/MrJP2LLkG/hF7GEQzaHo2aS7ho6ZLCOvf+6g==", + "engines": { + "node": ">=12.0.0" + } + }, "node_modules/prompts": { "version": "2.4.2", "resolved": "https://registry.npmjs.org/prompts/-/prompts-2.4.2.tgz", @@ -18270,6 +18320,11 @@ "resolved": "https://registry.npmjs.org/pg-connection-string/-/pg-connection-string-2.6.2.tgz", "integrity": "sha512-ch6OwaeaPYcova4kKZ15sbJ2hKb/VP48ZD2gE7i1J+L4MspCtBMAx8nMgz7bksc7IojCIIWuEhHibSMFH8m8oA==" }, + "pg-format": { + "version": "1.0.4", + "resolved": "https://registry.npmjs.org/pg-format/-/pg-format-1.0.4.tgz", + "integrity": "sha512-YyKEF78pEA6wwTAqOUaHIN/rWpfzzIuMh9KdAhc3rSLQ/7zkRFcCgYBAEGatDstLyZw4g0s9SNICmaTGnBVeyw==" + }, "pg-int8": { "version": "1.0.1", "resolved": "https://registry.npmjs.org/pg-int8/-/pg-int8-1.0.1.tgz", @@ -18286,6 +18341,17 @@ "resolved": "https://registry.npmjs.org/pg-protocol/-/pg-protocol-1.6.0.tgz", "integrity": "sha512-M+PDm637OY5WM307051+bsDia5Xej6d9IR4GwJse1qA1DIhiKlksvrneZOYQq42OM+spubpcNYEo2FcKQrDk+Q==" }, + "pg-pubsub": { + "version": "0.8.1", + "resolved": "https://registry.npmjs.org/pg-pubsub/-/pg-pubsub-0.8.1.tgz", + "integrity": "sha512-b/EHOwCrag4isghc4XgRipeAjfgyNg1DiL3Dwwh1Ojp91Lriltn5eg2nSWjBe4pzcFzhTM6HiB7LOG9NN1nx5g==", + "requires": { + "pg": "^8.7.3", + "pg-format": "^1.0.2", + "pony-cause": "^2.1.8", + "promised-retry": "^0.5.0" + } + }, "pg-types": { "version": "2.2.0", "resolved": "https://registry.npmjs.org/pg-types/-/pg-types-2.2.0.tgz", @@ -18410,6 +18476,11 @@ "integrity": "sha512-Nc3IT5yHzflTfbjgqWcCPpo7DaKy4FnpB0l/zCAW0Tc7jxAiuqSxHasntB3D7887LSrA93kDJ9IXovxJYxyLCA==", "dev": true }, + "pony-cause": { + "version": "2.1.11", + "resolved": "https://registry.npmjs.org/pony-cause/-/pony-cause-2.1.11.tgz", + "integrity": "sha512-M7LhCsdNbNgiLYiP4WjsfLUuFmCfnjdF6jKe2R9NKl4WFN+HZPGHJZ9lnLP7f9ZnKe3U9nuWD0szirmj+migUg==" + }, "postgres-array": { "version": "2.0.0", "resolved": "https://registry.npmjs.org/postgres-array/-/postgres-array-2.0.0.tgz", @@ -18504,6 +18575,21 @@ "retry": "^0.12.0" } }, + "promised-retry": { + "version": "0.5.0", + "resolved": "https://registry.npmjs.org/promised-retry/-/promised-retry-0.5.0.tgz", + "integrity": "sha512-jbYvN6UGE+/3E1g0JmgDPchUc+4VI4cBaPjdr2Lso22xfFqut2warEf6IhWuhPJKbJYVOQAyCt2Jx+01ORCItg==", + "requires": { + "pony-cause": "^1.1.1" + }, + "dependencies": { + "pony-cause": { + "version": "1.1.1", + "resolved": "https://registry.npmjs.org/pony-cause/-/pony-cause-1.1.1.tgz", + "integrity": "sha512-PxkIc/2ZpLiEzQXu5YRDOUgBlfGYBY8156HY5ZcRAwwonMk5W/MrJP2LLkG/hF7GEQzaHo2aS7ho6ZLCOvf+6g==" + } + } + }, "prompts": { "version": "2.4.2", "resolved": "https://registry.npmjs.org/prompts/-/prompts-2.4.2.tgz", diff --git a/package.json b/package.json index fe01d6c..f73f382 100644 --- a/package.json +++ b/package.json @@ -41,6 +41,7 @@ "merkle-lib": "^2.0.10", "node-telegram-bot-api": "^0.61.0", "pg": "^8.11.3", + "pg-pubsub": "^0.8.1", "reflect-metadata": "^0.1.13", "rpc-bitcoin": "^2.0.0", "rxjs": "^7.2.0", @@ -92,4 +93,4 @@ "coverageDirectory": "../coverage", "testEnvironment": "node" } -} \ No newline at end of file +} diff --git a/src/ORM/blocks/blocks.entity.ts b/src/ORM/blocks/blocks.entity.ts index ae1fe99..01f76c1 100644 --- a/src/ORM/blocks/blocks.entity.ts +++ b/src/ORM/blocks/blocks.entity.ts @@ -5,7 +5,7 @@ import { TrackedEntity } from '../utils/TrackedEntity.entity'; @Entity() export class BlocksEntity extends TrackedEntity { - @PrimaryGeneratedColumn() + @PrimaryGeneratedColumn({type: 'bigint'}) id: number; @Column() diff --git a/src/ORM/client-statistics/client-statistics.entity.ts b/src/ORM/client-statistics/client-statistics.entity.ts index a449505..e319e46 100644 --- a/src/ORM/client-statistics/client-statistics.entity.ts +++ b/src/ORM/client-statistics/client-statistics.entity.ts @@ -8,7 +8,7 @@ import { TrackedEntity } from '../utils/TrackedEntity.entity'; @Index(["clientId", "time"]) export class ClientStatisticsEntity extends TrackedEntity { - @PrimaryGeneratedColumn() + @PrimaryGeneratedColumn({type: 'bigint'}) id: number; @Column({ length: 62, type: 'varchar' }) diff --git a/src/ORM/home-graph/home-graph.entity.ts b/src/ORM/home-graph/home-graph.entity.ts index b62a6d7..908619d 100644 --- a/src/ORM/home-graph/home-graph.entity.ts +++ b/src/ORM/home-graph/home-graph.entity.ts @@ -3,7 +3,7 @@ import { Column, Entity, PrimaryGeneratedColumn } from 'typeorm'; @Entity() export class HomeGraphEntity { - @PrimaryGeneratedColumn() + @PrimaryGeneratedColumn({type: 'bigint'}) id: number; @Column({ type: 'bigint' }) diff --git a/src/ORM/rpc-block/rpc-block.service.ts b/src/ORM/rpc-block/rpc-block.service.ts index 75f1d54..28a7bec 100644 --- a/src/ORM/rpc-block/rpc-block.service.ts +++ b/src/ORM/rpc-block/rpc-block.service.ts @@ -12,7 +12,7 @@ export class RpcBlockService { ) { } - public getBlock(blockHeight: number) { + public getSavedBlockTemplate(blockHeight: number) { return this.rpcBlockRepository.findOne({ where: { blockHeight } }); diff --git a/src/ORM/telegram-subscriptions/telegram-subscriptions.entity.ts b/src/ORM/telegram-subscriptions/telegram-subscriptions.entity.ts index ead2294..011b758 100644 --- a/src/ORM/telegram-subscriptions/telegram-subscriptions.entity.ts +++ b/src/ORM/telegram-subscriptions/telegram-subscriptions.entity.ts @@ -5,7 +5,7 @@ import { TrackedEntity } from '../utils/TrackedEntity.entity'; @Entity() export class TelegramSubscriptionsEntity extends TrackedEntity { - @PrimaryGeneratedColumn() + @PrimaryGeneratedColumn({type: 'bigint'}) id: number; @Index() diff --git a/src/app.controller.ts b/src/app.controller.ts index ae7ca88..07464d3 100644 --- a/src/app.controller.ts +++ b/src/app.controller.ts @@ -71,7 +71,7 @@ export class AppController { const userAgents = await this.userAgentReportService.getReport(); const totalHashRate = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.totalHashRate), 0); const totalMiners = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.count), 0); - const blockHeight = (await firstValueFrom(this.bitcoinRpcService.newBlock$)).blocks; + const blockHeight = this.bitcoinRpcService.miningInfo.blocks; const blocksFound = await this.blocksService.getFoundBlocks(); const data = { @@ -90,8 +90,7 @@ export class AppController { @Get('network') public async network() { - const miningInfo = await firstValueFrom(this.bitcoinRpcService.newBlock$); - return miningInfo; + return this.bitcoinRpcService.miningInfo; } @Get('info/chart') diff --git a/src/services/bitcoin-rpc.service.ts b/src/services/bitcoin-rpc.service.ts index c924e48..5612b67 100644 --- a/src/services/bitcoin-rpc.service.ts +++ b/src/services/bitcoin-rpc.service.ts @@ -1,30 +1,38 @@ import { Injectable, OnModuleInit } from '@nestjs/common'; import { ConfigService } from '@nestjs/config'; import { RPCClient } from 'rpc-bitcoin'; -import { BehaviorSubject, filter, shareReplay } from 'rxjs'; +import { asyncScheduler, BehaviorSubject, delay, filter, from, interval, scheduled, shareReplay, startWith, Subject, switchMap } from 'rxjs'; import { RpcBlockService } from 'src/ORM/rpc-block/rpc-block.service'; import * as zmq from 'zeromq'; import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate'; import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo'; +import * as PGPubsub from 'pg-pubsub'; @Injectable() export class BitcoinRpcService implements OnModuleInit { - private blockHeight = 0; + private client: RPCClient; - private _newBlock$: BehaviorSubject = new BehaviorSubject(undefined); - public newBlock$ = this._newBlock$.pipe(filter(block => block != null), shareReplay({ refCount: true, bufferSize: 1 })); + private _newBlockTemplate$: BehaviorSubject = new BehaviorSubject(undefined); + private pubsubInstance: PGPubsub; + private resetTemplateInterval$ = new Subject(); + + public miningInfo: IMiningInfo; + public newBlockTemplate$ = this._newBlockTemplate$.pipe(filter(block => block != null), shareReplay({ refCount: true, bufferSize: 1 })); constructor( private readonly configService: ConfigService, private rpcBlockService: RpcBlockService ) { - } async onModuleInit() { + + this.pubsubInstance = new PGPubsub('postgres://' + this.configService.get('DB_USERNAME') + ':' + this.configService.get('DB_PASSWORD') + '@' + this.configService.get('DB_HOST') + ':' + this.configService.get('DB_PORT') + '/' + this.configService.get('DB_DATABASE')) + + const url = this.configService.get('BITCOIN_RPC_URL'); const user = this.configService.get('BITCOIN_RPC_USER'); const pass = this.configService.get('BITCOIN_RPC_PASSWORD'); @@ -38,12 +46,21 @@ export class BitcoinRpcService implements OnModuleInit { }, () => { console.error('Could not reach RPC host'); }); + + this.miningInfo = await this.getMiningInfo(); - if (this.configService.get('BITCOIN_ZMQ_HOST')) { + console.log(`MASTER? ${process.env.MASTER}`) + if (process.env.MASTER != 'true') { + this.pubsubInstance.addChannel('miningInfo', async (miningInfo: IMiningInfo) => { + console.log('PG Sub. new template'); + this.miningInfo = miningInfo; + const savedBlockTemplate = await this.rpcBlockService.getSavedBlockTemplate(miningInfo.blocks); + this._newBlockTemplate$.next(JSON.parse(savedBlockTemplate.data)); + }); + } else { console.log('Using ZMQ'); const sock = new zmq.Subscriber; - sock.connectTimeout = 1000; sock.events.on('connect', () => { console.log('ZMQ Connected'); @@ -56,70 +73,39 @@ export class BitcoinRpcService implements OnModuleInit { 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); + // Between new blocks we want refresh jobs with the latest transactions + this.resetTemplateInterval$.pipe( + startWith(null), + switchMap(() =>interval(60000)) + ).subscribe(async () =>{ + await this.getAndBroadcastLatestTemplate(); + }); + } + } private async listenForNewBlocks(sock: zmq.Subscriber) { for await (const [topic, msg] of sock) { console.log("New Block"); - await this.pollMiningInfo(); + this.miningInfo = await this.getMiningInfo(); + await this.getAndBroadcastLatestTemplate(); + + //Reset the block update interval + this.resetTemplateInterval$.next(); } } - public async pollMiningInfo() { - const miningInfo = await this.getMiningInfo(); - if (miningInfo != null && miningInfo.blocks > this.blockHeight) { - console.log("block height change"); - this._newBlock$.next(miningInfo); - this.blockHeight = miningInfo.blocks; - } - } - - private async waitForBlock(blockHeight: number): Promise { - while (true) { - await new Promise(r => setTimeout(r, 100)); - - const block = await this.rpcBlockService.getBlock(blockHeight); - if (block != null && block.data != null) { - console.log(`promise loop resolved, block height ${blockHeight}`); - return Promise.resolve(JSON.parse(block.data)); - } - console.log(`promise loop, block height ${blockHeight}`); - } - } - - public async getBlockTemplate(blockHeight: number): Promise { - let result: IBlockTemplate; - try { - const block = await this.rpcBlockService.getBlock(blockHeight); - const completeBlock = block?.data != null; - - - if(process.env.MASTER == 'true'){ - result = await this.loadBlockTemplate(blockHeight); - } - - if (completeBlock) { - return Promise.resolve(JSON.parse(block.data)); - } else{ - result = await this.waitForBlock(blockHeight); - } - - } catch (e) { - console.error('Error getblocktemplate:', e.message); - throw new Error('Error getblocktemplate'); - } - console.log(`getblocktemplate tx count: ${result.transactions.length}`); - return result; + public async getAndBroadcastLatestTemplate() { + const blockTemplate = await this.loadBlockTemplate(this.miningInfo.blocks); + this._newBlockTemplate$.next(blockTemplate); + await this.pubsubInstance.publish('miningInfo', this.miningInfo); } private async loadBlockTemplate(blockHeight: number) { - console.log(`Master fetching block ${blockHeight}`); + console.log(`Master fetching block template ${blockHeight}`); let blockTemplate: IBlockTemplate; while (blockTemplate == null) { @@ -132,11 +118,11 @@ export class BitcoinRpcService implements OnModuleInit { }); } - try{ + try { console.log(`Saving block ${blockHeight}`); await this.rpcBlockService.saveBlock(blockHeight, JSON.stringify(blockTemplate)); console.log('block saved'); - }catch(e){ + } catch (e) { console.error('Error saving block', e); } diff --git a/src/services/stratum-v1-jobs.service.ts b/src/services/stratum-v1-jobs.service.ts index 8c0c1dd..056a53e 100644 --- a/src/services/stratum-v1-jobs.service.ts +++ b/src/services/stratum-v1-jobs.service.ts @@ -23,49 +23,32 @@ export interface IJobTemplate { @Injectable() export class StratumV1JobsService { - private lastIntervalCount: number; - private skipNext: boolean = false; public newMiningJob$: Observable; - public latestJobId: number = 1; public latestJobTemplateId: number = 1; - public jobs: { [jobId: string]: MiningJob } = {}; - public blocks: { [id: number]: IJobTemplate } = {}; - // offset the interval so that all the cluster processes don't try and refresh at the same time. - private delay = process.env.NODE_APP_INSTANCE == null ? 0 : parseInt(process.env.NODE_APP_INSTANCE) * 5000; + private lastBlockHeight = 0; constructor( private readonly bitcoinRpcService: BitcoinRpcService ) { - this.newMiningJob$ = combineLatest([this.bitcoinRpcService.newBlock$, interval(60000).pipe(delay(this.delay), startWith(-1))]).pipe( - switchMap(([miningInfo, interval]) => { - return from(this.bitcoinRpcService.getBlockTemplate(miningInfo.blocks)).pipe(map((blockTemplate) => { - return { - blockTemplate, - interval - } - })) - }), - map(({ blockTemplate, interval }) => { + this.newMiningJob$ = this.bitcoinRpcService.newBlockTemplate$.pipe( + map((blockTemplate) => { + + console.log('Updating block template'); let clearJobs = false; - if (this.lastIntervalCount === interval) { + const currentBlockHeight = this.bitcoinRpcService.miningInfo.blocks; + + if(this.lastBlockHeight == 0 || this.lastBlockHeight != currentBlockHeight){ + console.log('New template is new block, clearing jobs'); clearJobs = true; - this.skipNext = true; - console.log('new block') + this.lastBlockHeight = currentBlockHeight; } - if (this.skipNext == true && clearJobs == false) { - this.skipNext = false; - return null; - } - - this.lastIntervalCount = interval; - const currentTime = Math.floor(new Date().getTime() / 1000); return { version: blockTemplate.version, @@ -130,6 +113,8 @@ export class StratumV1JobsService { }), shareReplay({ refCount: true, bufferSize: 1 }) ) + + this.newMiningJob$.subscribe(); } private calculateNetworkDifficulty(nBits: number) {