131 Commits
Author SHA1 Message Date
Benjamin WilsonandWantClue 9bf3dde862 validateHeaderCompliance (#145)
Co-authored-by: WantClue <wantclue@users.noreply.github.com>
2026-01-24 14:53:05 -05:00
Ben 8b2bea702a enable stratum TLS 2025-11-29 11:45:44 -05:00
Ben 63a1618f26 killDeadClients 2025-11-21 20:56:27 -05:00
Ben 1b0f3f52d8 index 2025-11-21 20:52:52 -05:00
Ben 78c00fd25c killDeadClients 2025-11-21 20:34:00 -05:00
Ben dd45c60eff killDeadClients 2025-11-21 20:01:55 -05:00
Ben 0481dfda6e dead clients 2025-11-21 18:56:29 -05:00
Ben 6f6abe5e2d batch kill dead clients 2025-11-21 18:42:00 -05:00
Ben ff84d8284f watch cert and reload 2025-08-22 19:01:46 -04:00
Ben cdf6eb807b other filtering 2025-04-29 15:31:28 -04:00
Benjamin Wilson d9c89347e1 remove logging 2025-04-27 13:13:05 -04:00
Ben c1d8341c22 Prevent duplicate share submissions - this is why we can't have nice things 2025-04-27 13:07:24 -04:00
Benjamin Wilson f39c0dc059 Increase miner compatibility by allowing, but stripping unknown message values rather than rejecting as validation errors 2025-04-19 13:50:50 -04:00
Benjamin Wilson 3dab14e3b7 increase timeout to 15min 2025-03-03 20:29:31 -05:00
Benjamin Wilson 5060345809 better way to calculate hashrate 2025-02-27 16:22:48 -05:00
Benjamin Wilson b44f09848e don't calculate hashrate for 1 min 2025-02-27 12:15:20 -05:00
Benjamin Wilson 9b5402f1b0 refactor hashrate calculation 2025-02-27 12:10:43 -05:00
Benjamin Wilson 17881f70d6 heartbeat 2025-02-26 20:20:37 -05:00
Benjamin Wilson e9949f4038 utc string 2025-02-26 20:15:24 -05:00
Benjamin Wilson 7c3b31f7b3 timezone 2025-02-26 19:59:33 -05:00
Benjamin Wilson 78deb3a375 logging 2025-02-26 17:52:32 -05:00
Benjamin Wilson c5c8810da6 heartbeat updates 2025-02-26 17:48:27 -05:00
Benjamin Wilson 8584c442dc bug 2025-02-26 17:07:04 -05:00
Benjamin Wilson fd24807c55 cleanup 2025-02-26 09:50:42 -05:00
Benjamin Wilson bd5333a725 fix data type 2025-02-26 09:16:46 -05:00
Benjamin Wilson 419eb09853 logging 2025-02-26 00:32:43 -05:00
Benjamin Wilson 74f5783dc3 logging 2025-02-26 00:24:34 -05:00
Benjamin Wilson 34a27d0f2f update all stats 2025-02-26 00:18:25 -05:00
Benjamin Wilson d8e087f720 server start setTimeout 2025-02-26 00:08:20 -05:00
Benjamin Wilson 34eb0ca0f2 logging 2025-02-26 00:06:25 -05:00
Benjamin Wilson a33d7e07b3 logging 2025-02-25 23:59:29 -05:00
Benjamin Wilson f8d1c61c1c bug 2025-02-25 23:40:50 -05:00
Benjamin Wilson 326a5a7b33 typeo 2025-02-25 19:27:34 -05:00
Benjamin Wilson 903aa2628c query 2025-02-25 19:17:13 -05:00
Benjamin Wilson 19f49b74e6 sql camel case 2025-02-25 19:11:10 -05:00
Benjamin Wilson 864867ffd4 remove socket server delay 2025-02-25 19:06:33 -05:00
Benjamin Wilson f465b837f5 bulk client stat updates 2025-02-25 19:02:38 -05:00
Remco Ros 637d30f1a5 patch rpc-bitcoin to be compatible with bitcoin core 28.0 (#66) 2025-02-24 19:37:47 -05:00
Benjamin Wilson df87214de9 Online device list compaction 2025-02-23 12:57:38 -05:00
Benjamin Wilson 66ce822faa Fix kill dead clients tx wantclue 2025-02-22 17:47:46 -05:00
Benjamin Wilson 78ef48b805 logging 2025-02-10 00:07:14 -05:00
Benjamin Wilson b545e74d4f delete old templates and jobs 2025-02-10 00:01:55 -05:00
Benjamin Wilson 5702e20dc0 reduce the user agent list 2025-02-09 19:38:08 -05:00
Benjamin Wilson 74dd5ba66f Use pg pub/sub for stateless app notification 2025-02-09 17:04:31 -05:00
Benjamin Wilson 674b02ca9d socket cleanup 2025-02-03 15:39:03 -05:00
Benjamin Wilson 40ee610600 logs 2025-02-01 12:36:07 -05:00
Benjamin Wilson 7e0de263f2 Revert "Fix home graph hashrate calculation (around restarts)"
This reverts commit 220a1d1558.
2025-02-01 12:35:55 -05:00
Benjamin Wilson 65232da901 block save => upsert 2025-02-01 11:57:34 -05:00
Benjamin Wilson 220a1d1558 Fix home graph hashrate calculation (around restarts) 2025-01-02 23:34:43 -05:00
Benjamin Wilson ae7ed021a3 logging 2025-01-02 11:15:46 -05:00
Benjamin Wilson 056fe7af4f logging 2025-01-01 23:06:19 -05:00
Benjamin Wilson 912ae9fe85 tweak timeout and starting diff 2024-12-22 10:03:03 -05:00
Benjamin Wilson d513dab932 user agent cleanup 2024-10-16 13:53:49 -04:00
Benjamin Wilson 3a1f3bba10 Ensure mintime is considered for jobs 2024-09-27 13:00:52 -04:00
Benjamin Wilson 8fdf1e839e save block 2024-09-16 19:48:46 -04:00
Benjamin Wilson 76d92d84b4 debug 2024-09-16 19:43:24 -04:00
Benjamin Wilson 4928f34daf debug 2024-09-16 19:39:56 -04:00
Benjamin Wilson 88e7f5f3fb Saving block error 2024-09-16 19:35:13 -04:00
Benjamin Wilson d97843af3d getBlockTemplate 2024-09-16 19:31:40 -04:00
Benjamin Wilson c1ab6a152e loadBlockTemplate 2024-09-16 19:28:00 -04:00
Benjamin Wilson 717a5c7878 getBlockTemplate fix 2024-09-16 19:24:34 -04:00
Benjamin Wilson 6038af8448 fix getBlockTemplate 2024-09-16 19:16:38 -04:00
Benjamin Wilson d6863fad28 timing adjustments 2024-09-16 18:32:43 -04:00
Benjamin Wilson b9b5fe2cf5 cluster 2024-09-09 20:27:17 -04:00
Benjamin Wilson 9db0afd9df removed NODE_APP_INSTANCE 0 dependency 2024-09-09 20:18:37 -04:00
Benjamin Wilson fb8715c23d multi port 2024-09-07 23:56:09 -04:00
Benjamin Wilson ff6d446c12 validation handling 2024-08-22 15:18:40 -04:00
Benjamin Wilson 02d718934d initialize hashrate to 0 2024-07-17 17:24:04 -04:00
Benjamin Wilson d40742ef71 bug 2024-07-15 09:16:55 -04:00
Benjamin Wilson 38ecbdc7ec logging 2024-07-14 11:55:49 -04:00
Benjamin Wilson 8905409880 null jobTemplate 2024-07-10 22:40:02 -04:00
Benjamin Wilson 2bfbc2d79a logging 2024-07-10 21:51:46 -04:00
Benjamin Wilson 2ee82a269c loadBlockTemplate 2024-07-05 20:22:06 -04:00
Benjamin Wilson fe4d58a60d StratumBaseMessage validation 2024-07-05 20:21:47 -04:00
Benjamin Wilson b858a594b8 open stratum base message to accept number or string 2024-07-05 08:19:54 -04:00
Benjamin Wilson 3e4cb4ff5d always clear jobs on diff change 2024-07-05 08:13:55 -04:00
Benjamin Wilson 14f95a7fa4 tweaks 2024-07-04 21:03:51 -04:00
Ben Wilson 522ef126d9 updated at, hashrate calculation 2024-07-04 20:42:37 -04:00
Benjamin Wilson 0c34f23284 ensure block template is fresh 2024-07-04 20:06:48 -04:00
Benjamin Wilson e52d14b0a9 handle partial socket data 2024-07-04 17:19:58 -04:00
Benjamin Wilson 3727773a2e getBlockTemplate improvements 2024-07-04 16:46:57 -04:00
Benjamin Wilson ad3035b908 ensure getblocktemplate always returns a value 2024-07-04 16:46:44 -04:00
Benjamin Wilson 3e279ee662 remove dead code 2024-07-04 16:43:02 -04:00
Ben Wilson deeb40967e decrease submission rate 2024-04-30 16:45:08 -04:00
Ben Wilson f1221651f6 diff fix 2024-03-29 12:17:30 -04:00
Ben Wilson 32a604e745 allow client to suggest diff in password field "d=1234" 2024-03-27 17:02:11 -04:00
Georges 6442bb45e0 stratum: set_difficulty should always be a notification (#32) 2024-02-23 12:54:16 -05:00
Ben Wilson f2930681d3 optional password 2024-02-23 12:54:16 -05:00
Ben Wilson 3d3d1a55a8 share stats 2024-02-18 14:38:45 -05:00
Ben Wilson 01adb0a997 ClientStatistics updatedAt 2024-02-18 12:21:52 -05:00
Ben Wilson 5ead6c605b remove server start delay 2024-02-18 10:52:28 -05:00
Ben Wilson 4fd8d9ceb6 change coinbase tag 2024-02-17 20:31:23 -05:00
Ben Wilson 9d4c3ad35e lastSave 2024-02-17 20:12:37 -05:00
Ben Wilson aab69b65fc clientid 2024-02-17 19:01:24 -05:00
Ben Wilson 9895e80a01 share updates 2024-02-17 18:40:34 -05:00
Ben Wilson 9617ec5685 statistics saving 2024-02-17 18:25:14 -05:00
Ben Wilson aefb180fb4 fix pool info 2024-02-17 16:27:37 -05:00
Ben Wilson 87508bfb69 user agent report 2024-02-16 23:49:47 -05:00
Ben Wilson 39e537a00c bug 2024-02-16 17:18:21 -05:00
Ben Wilson 9ecfc68155 bug 2024-02-16 17:17:19 -05:00
Ben Wilson 3e2de62259 db refactor 2024-02-16 17:12:36 -05:00
Ben Wilson 8e977c21a1 Merge remote-tracking branch 'origin/master' into postgresql 2024-02-14 08:38:05 -05:00
Ben Wilson 5fd9f98807 index, dead code 2024-01-22 17:32:34 -05:00
Ben Wilson 6bc42defef separate home graph into it's own table 2024-01-16 19:37:29 -05:00
Ben Wilson bbc14e2272 bug 2024-01-15 20:04:26 -05:00
Ben Wilson fc767fa770 reduce heartbeat updates 2024-01-14 14:10:47 -05:00
Ben Wilson 40e32d6e1a stratum jobs array to dictionary 2023-12-15 00:56:22 -05:00
Ben Wilson df6ec07c58 batch save share statistics 2023-12-15 00:25:01 -05:00
Ben Wilson df245f396e parameterized typeorm 2023-12-14 23:54:45 -05:00
Ben Wilson 42fc863c7c column names 2023-12-14 23:47:51 -05:00
Ben Wilson e46bb24155 raw query 2023-12-14 23:44:41 -05:00
Ben Wilson cb4a173733 add back index with time 2023-12-14 22:45:12 -05:00
Ben Wilson 83c769849d removed hash from client stats table, connection pool size 2023-12-14 22:30:55 -05:00
Ben Wilson ee5fea1df5 wrap save in a transaction 2023-12-10 17:23:08 -05:00
Ben Wilson 06acad6770 example 2023-12-08 00:12:25 -05:00
Ben Wilson f7b22bb310 delete old stats sooner 2023-12-07 09:57:20 -05:00
Ben Wilson 6cf6a007d9 remove read UNCOMMITTED 2023-12-07 09:55:36 -05:00
Ben Wilson 5b71632f8a postgres doesn't support read UNCOMMITTED 2023-12-07 09:55:07 -05:00
Ben Wilson c5ff8521ac Merge branch 'postgresql' of https://github.com/benjamin-wilson/public-pool into postgresql 2023-12-07 09:53:40 -05:00
Ben Wilson 53ee510a1e fix chart, caching 2023-12-07 09:53:16 -05:00
Ben Wilson 7be554a358 fix chart, caching 2023-12-07 09:40:32 -05:00
Ben Wilson eb5d434157 few optimizations 2023-12-07 09:07:37 -05:00
Ben Wilson 893c7c2003 revert 2023-12-07 08:33:43 -05:00
Ben Wilson 4736eadf7b row names 2023-12-07 08:19:51 -05:00
Ben Wilson 1089e6210c deadlock? 2023-12-07 08:16:48 -05:00
Ben Wilson 017978d250 userAgent 2023-12-04 08:52:49 -05:00
Ben Wilson 5a17528b8f bug 2023-12-04 08:33:47 -05:00
Ben Wilson 0b147ea5db fix numbers returned from query as string 2023-12-04 08:26:47 -05:00
Ben Wilson bffcf6c74c config 2023-12-04 01:07:45 -05:00
Ben Wilson 96a38f9d9b bug 2023-12-04 01:06:53 -05:00
Ben Wilson 39c4c033e2 postgres init 2023-12-04 00:34:33 -05:00
47 changed files with 2191 additions and 913 deletions
+11 -6
View File
@@ -15,7 +15,10 @@ BITCOIN_RPC_TIMEOUT=10000
# BITCOIN_ZMQ_HOST="tcp://192.168.1.100:3000" # BITCOIN_ZMQ_HOST="tcp://192.168.1.100:3000"
API_PORT=3334 API_PORT=3334
STRATUM_PORT=3333 STRATUM_PORTS=3333,3332,3331,3330
STRATUM_SECURE=true
SECURE_STRATUM_PORTS=4333,4332,4331,4330
#optional telegram bot #optional telegram bot
#TELEGRAM_BOT_TOKEN= #TELEGRAM_BOT_TOKEN=
@@ -32,9 +35,11 @@ NETWORK=mainnet
API_SECURE=false API_SECURE=false
ENABLE_SOLO=false #postgresql
ENABLE_PROXY=true DB_HOST=
DB_PORT=
DB_USERNAME=
DB_PASSWORD=
DB_DATABASE=
BRAIINS_ACCESS_TOKEN= PRODUCTION=false
PROXY_PORT=3333
+25
View File
@@ -0,0 +1,25 @@
module.exports = {
apps: [
// Master instance
{
name: 'master',
script: './dist/main.js',
instances: 1,
env: {
MASTER: 'true',
},
time: true
},
// Worker instances
{
name: 'workers',
script: './dist/main.js',
instances: 2,
exec_mode: "cluster",
env: {
MASTER: 'false',
},
time: true
},
],
};
-7
View File
@@ -25,10 +25,3 @@ DEV_FEE_ADDRESS=
NETWORK=mainnet NETWORK=mainnet
API_SECURE=false API_SECURE=false
ENABLE_SOLO=true
ENABLE_PROXY=false
BRAIINS_ACCESS_TOKEN=
PROXY_PORT=3333
-7
View File
@@ -25,10 +25,3 @@ DEV_FEE_ADDRESS=
NETWORK=regtest NETWORK=regtest
API_SECURE=false API_SECURE=false
ENABLE_SOLO=true
ENABLE_PROXY=false
BRAIINS_ACCESS_TOKEN=
PROXY_PORT=3333
+1 -8
View File
@@ -24,11 +24,4 @@ DEV_FEE_ADDRESS=
# mainnet | testnet # mainnet | testnet
NETWORK=testnet NETWORK=testnet
API_SECURE=false API_SECURE=false
ENABLE_SOLO=true
ENABLE_PROXY=false
BRAIINS_ACCESS_TOKEN=
PROXY_PORT=3333
+857 -48
View File
File diff suppressed because it is too large Load Diff
+5 -1
View File
@@ -8,6 +8,7 @@
"scripts": { "scripts": {
"build": "nest build", "build": "nest build",
"format": "prettier --write \"src/**/*.ts\" \"test/**/*.ts\"", "format": "prettier --write \"src/**/*.ts\" \"test/**/*.ts\"",
"postinstall": "patch-package",
"start": "nest start", "start": "nest start",
"start:dev": "nest start --watch", "start:dev": "nest start --watch",
"start:debug": "nest start --debug --watch", "start:debug": "nest start --debug --watch",
@@ -40,6 +41,8 @@
"discord.js": "^14.11.0", "discord.js": "^14.11.0",
"merkle-lib": "^2.0.10", "merkle-lib": "^2.0.10",
"node-telegram-bot-api": "^0.61.0", "node-telegram-bot-api": "^0.61.0",
"pg": "^8.11.3",
"pg-pubsub": "^0.8.1",
"reflect-metadata": "^0.1.13", "reflect-metadata": "^0.1.13",
"rpc-bitcoin": "^2.0.0", "rpc-bitcoin": "^2.0.0",
"rxjs": "^7.2.0", "rxjs": "^7.2.0",
@@ -65,6 +68,7 @@
"eslint-config-prettier": "^8.3.0", "eslint-config-prettier": "^8.3.0",
"eslint-plugin-prettier": "^4.0.0", "eslint-plugin-prettier": "^4.0.0",
"jest": "29.5.0", "jest": "29.5.0",
"patch-package": "8.0.0",
"prettier": "^2.3.2", "prettier": "^2.3.2",
"source-map-support": "^0.5.20", "source-map-support": "^0.5.20",
"supertest": "^6.1.3", "supertest": "^6.1.3",
@@ -91,4 +95,4 @@
"coverageDirectory": "../coverage", "coverageDirectory": "../coverage",
"testEnvironment": "node" "testEnvironment": "node"
} }
} }
+26
View File
@@ -0,0 +1,26 @@
diff --git a/node_modules/rpc-bitcoin/build/src/rpc.d.ts b/node_modules/rpc-bitcoin/build/src/rpc.d.ts
index d25d732..ce4eb22 100644
--- a/node_modules/rpc-bitcoin/build/src/rpc.d.ts
+++ b/node_modules/rpc-bitcoin/build/src/rpc.d.ts
@@ -6,7 +6,7 @@ export declare type RPCIniOptions = RESTIniOptions & {
fullResponse?: boolean;
};
export declare type JSONRPC = {
- jsonrpc?: string | number;
+ jsonrpc?: string;
id?: string | number;
method: string;
params?: object;
diff --git a/node_modules/rpc-bitcoin/build/src/rpc.js b/node_modules/rpc-bitcoin/build/src/rpc.js
index 7fec7aa..3bee5fc 100644
--- a/node_modules/rpc-bitcoin/build/src/rpc.js
+++ b/node_modules/rpc-bitcoin/build/src/rpc.js
@@ -12,7 +12,7 @@ class RPCClient extends rest_1.RESTClient {
}
async rpc(method, params = {}, wallet) {
const uri = typeof wallet === "undefined" ? "/" : "wallet/" + wallet;
- const body = { method, params, jsonrpc: 1.0, id: "rpc-bitcoin" };
+ const body = { method, params, jsonrpc: "1.0", id: "rpc-bitcoin" };
try {
const response = await this.batch(body, uri);
return this.fullResponse ? response : response.result;
+15
View File
@@ -0,0 +1,15 @@
import { MigrationInterface, QueryRunner } from 'typeorm';
export class UniqueNonceIndex implements MigrationInterface {
public async up(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(
`CREATE UNIQUE INDEX "IDX_unique_nonce" ON "client_entity" ("sessionId") WHERE "deletedAt" IS NOT NULL`
);
}
public async down(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`DROP INDEX "IDX_unique_nonce"`);
}
}
@@ -0,0 +1,15 @@
import { Global, Module } from '@nestjs/common';
import { TypeOrmModule } from '@nestjs/typeorm';
import { UserAgentReportService } from './user-agent-report.service';
import { UserAgentReportView } from './user-agent-report.view';
@Global()
@Module({
imports: [TypeOrmModule.forFeature([UserAgentReportView])],
providers: [UserAgentReportService],
exports: [TypeOrmModule, UserAgentReportService],
})
export class UserAgentReportModule { }
@@ -0,0 +1,35 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { UserAgentReportView } from './user-agent-report.view';
@Injectable()
export class UserAgentReportService {
constructor(
@InjectRepository(UserAgentReportView)
private userAgentReport: Repository<UserAgentReportView>,
) {
}
public async getReport() {
return await this.userAgentReport.find();
}
public async refreshReport() {
try {
return await this.userAgentReport.query(`
COMMIT;
BEGIN TRANSACTION ISOLATION LEVEL REPEATABLE READ;
REFRESH MATERIALIZED VIEW user_agent_report_view;
COMMIT;
`);
} catch (e) {
console.log(e)
}
}
}
@@ -0,0 +1,27 @@
import { DataSource, ViewColumn, ViewEntity } from 'typeorm';
import { ClientEntity } from '../../client/client.entity';
@ViewEntity({
materialized: true,
expression: (dataSource: DataSource) =>
dataSource
.createQueryBuilder()
.select('client.userAgent as "userAgent"')
.addSelect('COUNT(client.userAgent)', 'count')
.addSelect('MAX(client.bestDifficulty)', 'bestDifficulty')
.addSelect('SUM(client.hashRate)', 'totalHashRate')
.from(ClientEntity, 'client')
.groupBy('client.userAgent')
.orderBy('"totalHashRate"', 'DESC')
})
export class UserAgentReportView {
@ViewColumn()
userAgent: string;
@ViewColumn()
count: string;
@ViewColumn()
bestDifficulty: number;
@ViewColumn()
totalHashRate: string;
}
@@ -11,9 +11,12 @@ export class AddressSettingsEntity extends TrackedEntity {
@Column({ default: 0 }) @Column({ default: 0 })
shares: number; shares: number;
@Column({ type: 'real', default: 0 }) @Column({ type: 'decimal', default: 0 })
bestDifficulty: number; bestDifficulty: number;
@Column({ nullable: true })
bestDifficultyUserAgent: string;
@Column({ nullable: true }) @Column({ nullable: true })
miscCoinbaseScriptData: string; miscCoinbaseScriptData: string;
@@ -17,28 +17,33 @@ export class AddressSettingsService {
public async getSettings(address: string, createIfNotFound: boolean) { public async getSettings(address: string, createIfNotFound: boolean) {
const settings = await this.addressSettingsRepository.findOne({ where: { address } }); const settings = await this.addressSettingsRepository.findOne({ where: { address } });
if (createIfNotFound == true && settings == null) { if (createIfNotFound == true && settings == null) {
return await this.createNew(address); // It's possible to have a race condition here so if we get a PK violation, fetch it
try {
return await this.createNew(address);
} catch (e) {
return await this.addressSettingsRepository.findOne({ where: { address } });
}
} }
return settings; return settings;
} }
public async updateBestDifficulty(address: string, bestDifficulty: number) { public async updateBestDifficulty(address: string, bestDifficulty: number, bestDifficultyUserAgent: string) {
return await this.addressSettingsRepository.update({ address }, { bestDifficulty }); return await this.addressSettingsRepository.update({ address }, { bestDifficulty, bestDifficultyUserAgent });
} }
public async createNew(address: string) { public async createNew(address: string) {
return await this.addressSettingsRepository.save({ address }); return await this.addressSettingsRepository.save({ address });
} }
public async addShares(address: string, shares: number) { // public async addShares(address: string, shares: number) {
return await this.addressSettingsRepository.createQueryBuilder() // return await this.addressSettingsRepository.createQueryBuilder()
.update(AddressSettingsEntity) // .update(AddressSettingsEntity)
.set({ // .set({
shares: () => `"shares" + ${shares}` // Use the actual value of shares here // shares: () => `"shares" + ${shares}` // Use the actual value of shares here
}) // })
.where('address = :address', { address }) // .where('address = :address', { address })
.execute(); // .execute();
} // }
public async resetBestDifficultyAndShares() { public async resetBestDifficultyAndShares() {
return await this.addressSettingsRepository.update({}, { return await this.addressSettingsRepository.update({}, {
@@ -46,4 +51,12 @@ export class AddressSettingsService {
bestDifficulty: 0 bestDifficulty: 0
}); });
} }
public async getHighScores() {
return await this.addressSettingsRepository.createQueryBuilder()
.select('"updatedAt", "bestDifficulty", "bestDifficultyUserAgent"')
.orderBy('"bestDifficulty"', 'DESC')
.limit(10)
.execute();
}
} }
+1 -1
View File
@@ -5,7 +5,7 @@ import { TrackedEntity } from '../utils/TrackedEntity.entity';
@Entity() @Entity()
export class BlocksEntity extends TrackedEntity { export class BlocksEntity extends TrackedEntity {
@PrimaryGeneratedColumn() @PrimaryGeneratedColumn({type: 'bigint'})
id: number; id: number;
@Column() @Column()
@@ -1,15 +1,14 @@
import { Column, Entity, Index, PrimaryGeneratedColumn } from 'typeorm'; import { Column, Entity, Index, ManyToOne, PrimaryGeneratedColumn } from 'typeorm';
import { ClientEntity } from '../client/client.entity';
import { TrackedEntity } from '../utils/TrackedEntity.entity'; import { TrackedEntity } from '../utils/TrackedEntity.entity';
@Entity() @Entity()
//Index for getHashRateForSession
@Index(["address", "clientName", "sessionId"])
//Index for statistics save //Index for statistics save
@Index(["address", "clientName", "sessionId", "time"]) @Index(["clientId", "time"])
export class ClientStatisticsEntity extends TrackedEntity { export class ClientStatisticsEntity extends TrackedEntity {
@PrimaryGeneratedColumn() @PrimaryGeneratedColumn({type: 'bigint'})
id: number; id: number;
@Column({ length: 62, type: 'varchar' }) @Column({ length: 62, type: 'varchar' })
@@ -22,14 +21,27 @@ export class ClientStatisticsEntity extends TrackedEntity {
sessionId: string; sessionId: string;
@Index() @Index()
@Column({ type: 'integer' }) @Column({ type: 'bigint' })
time: number; time: number;
@Column({ type: 'real' }) @Column({ type: 'decimal' })
shares: number; shares: number;
@Column({ default: 0, type: 'integer' }) @Column({ default: 0, type: 'bigint' })
acceptedCount: number; acceptedCount: number;
@ManyToOne(
() => ClientEntity,
clientEntity => clientEntity.statistics,
{ nullable: false, }
)
client: ClientEntity;
@Index()
@Column({ name: 'clientId' })
public clientId: string;
} }
@@ -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 { ClientStatisticsEntity } from './client-statistics.entity'; import { ClientStatisticsEntity } from './client-statistics.entity';
@@ -8,38 +8,90 @@ import { ClientStatisticsEntity } from './client-statistics.entity';
@Injectable() @Injectable()
export class ClientStatisticsService { export class ClientStatisticsService {
private bulkAsyncUpdates: {
[key: string]: Partial<ClientStatisticsEntity>
} = {};
constructor( constructor(
@InjectDataSource()
private dataSource: DataSource,
@InjectRepository(ClientStatisticsEntity) @InjectRepository(ClientStatisticsEntity)
private clientStatisticsRepository: Repository<ClientStatisticsEntity>, private clientStatisticsRepository: Repository<ClientStatisticsEntity>,
) { ) {
} }
public async save(clientStatistic: Partial<ClientStatisticsEntity>) { // public async update(clientStatistic: Partial<ClientStatisticsEntity>) {
// Attempt to update the existing record
const updateResult = await this.clientStatisticsRepository
.createQueryBuilder()
.update(ClientStatisticsEntity)
.set({
shares: () => `"shares" + :sharesIncrement`,
acceptedCount: () => `"acceptedCount" + 1`
})
.where('address = :address AND clientName = :clientName AND sessionId = :sessionId AND time = :time', {
address: clientStatistic.address,
clientName: clientStatistic.clientName,
sessionId: clientStatistic.sessionId,
time: clientStatistic.time,
sharesIncrement: clientStatistic.shares
})
.execute();
// Check if the update affected any rows // await this.clientStatisticsRepository.update({ clientId: clientStatistic.clientId, time: clientStatistic.time },
if (updateResult.affected === 0) { // {
// If no rows were updated, insert a new record // shares: clientStatistic.shares,
await this.clientStatisticsRepository.insert(clientStatistic); // acceptedCount: clientStatistic.acceptedCount,
// updatedAt: new Date()
// });
// }
public updateBulkAsync(clientStatistic: Partial<ClientStatisticsEntity>) {
const key = clientStatistic.clientId + clientStatistic.time.toString();
if(this.bulkAsyncUpdates[key] != null){
this.bulkAsyncUpdates[key].shares = clientStatistic.shares;
this.bulkAsyncUpdates[key].acceptedCount = clientStatistic.acceptedCount;
return;
} }
this.bulkAsyncUpdates[clientStatistic.clientId + clientStatistic.time.toString()] = clientStatistic;
}
public async doBulkAsyncUpdate(){
if(Object.keys(this.bulkAsyncUpdates).length < 1){
console.log('No client stats to update.')
return;
}
// Step 1: Prepare data for bulk update
const values = Object.entries(this.bulkAsyncUpdates).map(([key, value]) => {
return `('${value.clientId}', ${value.time}, ${value.shares}, ${value.acceptedCount}, NOW())`
}).join(',');
const query = `
DO $$
BEGIN
CREATE TEMP TABLE temp_stats (
"clientId" UUID,
time BIGINT,
shares NUMERIC,
"acceptedCount" INT,
"updatedAt" TIMESTAMP
) ON COMMIT DROP;
INSERT INTO temp_stats ("clientId", time, shares, "acceptedCount", "updatedAt")
VALUES ${values};
UPDATE "client_statistics_entity" cse
SET shares = ts.shares,
"acceptedCount" = ts."acceptedCount",
"updatedAt" = ts."updatedAt"
FROM temp_stats ts
WHERE cse."clientId" = ts."clientId" AND cse.time = ts.time;
END;
$$;
`;
try {
await this.clientStatisticsRepository.query(query);
//console.log(`Bulk updated ${Object.keys(this.bulkAsyncUpdates).length} statistics`)
} catch (error) {
console.error('Bulk update failed:', error.message, query);
throw error;
}
this.bulkAsyncUpdates = {};
}
public async insert(clientStatistic: Partial<ClientStatisticsEntity>) {
await this.clientStatisticsRepository.insert(clientStatistic);
} }
public async deleteOldStatistics() { public async deleteOldStatistics() {
@@ -53,57 +105,7 @@ export class ClientStatisticsService {
.execute(); .execute();
} }
public async getChartDataForSite() {
var yesterday = new Date(new Date().getTime() - (24 * 60 * 60 * 1000));
const query = `
SELECT
time AS label,
ROUND(((SUM(shares) * 4294967296) / 600)) AS data
FROM
client_statistics_entity AS entry
WHERE
entry.time > ${yesterday.getTime()}
GROUP BY
time
ORDER BY
time
LIMIT 144;
`;
const result: any[] = await this.clientStatisticsRepository.query(query);
return result.map(res => {
res.label = new Date(res.label).toISOString();
return res;
}).slice(0, result.length - 1)
}
// public async getHashRateForAddress(address: string) {
// const oneHour = new Date(new Date().getTime() - (60 * 60 * 1000));
// const query = `
// SELECT
// SUM(entry.shares) AS difficultySum
// FROM
// client_statistics_entity AS entry
// WHERE
// entry.address = ? AND entry.time > ${oneHour}
// `;
// const result = await this.clientStatisticsRepository.query(query, [address]);
// const difficultySum = result[0].difficultySum;
// return (difficultySum * 4294967296) / (600);
// }
public async getChartDataForAddress(address: string) { public async getChartDataForAddress(address: string) {
@@ -111,12 +113,12 @@ export class ClientStatisticsService {
const query = ` const query = `
SELECT SELECT
time label, time AS label,
(SUM(shares) * 4294967296) / 600 AS data (SUM(shares) * 4294967296) / 600 AS data
FROM FROM
client_statistics_entity AS entry client_statistics_entity AS entry
WHERE WHERE
entry.address = ? AND entry.time > ${yesterday.getTime()} entry.address = $1 AND entry.time > $2
GROUP BY GROUP BY
time time
ORDER BY ORDER BY
@@ -125,10 +127,10 @@ export class ClientStatisticsService {
`; `;
const result = await this.clientStatisticsRepository.query(query, [address]); const result = await this.clientStatisticsRepository.query(query, [address, yesterday.getTime()]);
return result.map(res => { return result.map(res => {
res.label = new Date(res.label).toISOString(); res.label = new Date(parseInt(res.label)).toISOString();
return res; return res;
}).slice(0, result.length - 1); }).slice(0, result.length - 1);
@@ -146,7 +148,7 @@ export class ClientStatisticsService {
FROM FROM
client_statistics_entity AS entry client_statistics_entity AS entry
WHERE WHERE
entry.address = ? AND entry.clientName = ? AND entry.time > ${oneHour.getTime()} entry.address = $1 AND entry.clientName = $2 AND entry.time > ${oneHour.getTime()}
`; `;
const result = await this.clientStatisticsRepository.query(query, [address, clientName]); const result = await this.clientStatisticsRepository.query(query, [address, clientName]);
@@ -163,12 +165,12 @@ export class ClientStatisticsService {
const query = ` const query = `
SELECT SELECT
time label, time AS label,
(SUM(shares) * 4294967296) / 600 AS data (SUM(shares) * 4294967296) / 600 AS data
FROM FROM
client_statistics_entity AS entry client_statistics_entity AS entry
WHERE WHERE
entry.address = ? AND entry.clientName = ? AND entry.time > ${yesterday.getTime()} entry.address = $1 AND entry."clientName" = $2 AND entry.time > ${yesterday.getTime()}
GROUP BY GROUP BY
time time
ORDER BY ORDER BY
@@ -179,7 +181,7 @@ export class ClientStatisticsService {
const result = await this.clientStatisticsRepository.query(query, [address, clientName]); const result = await this.clientStatisticsRepository.query(query, [address, clientName]);
return result.map(res => { return result.map(res => {
res.label = new Date(res.label).toISOString(); res.label = new Date(parseInt(res.label)).toISOString();
return res; return res;
}).slice(0, result.length - 1); }).slice(0, result.length - 1);
@@ -187,57 +189,59 @@ export class ClientStatisticsService {
} }
public async getHashRateForSession(address: string, clientName: string, sessionId: string) { // public async getHashRateForSession(clientId: string) {
const query = ` // const query = `
SELECT // SELECT
createdAt, // "createdAt",
updatedAt, // "updatedAt",
shares // shares
FROM // FROM
client_statistics_entity AS entry // client_statistics_entity AS entry
WHERE // WHERE
entry.address = ? AND entry.clientName = ? AND entry.sessionId = ? // entry."clientId" = $1
ORDER BY time DESC // ORDER BY time DESC
LIMIT 2; // LIMIT 2;
`; // `;
const result = await this.clientStatisticsRepository.query(query, [address, clientName, sessionId]); // const result = await this.clientStatisticsRepository.query(query, [clientId]);
if (result.length < 1) { // if (result.length < 1) {
return 0; // return 0;
} // }
const latestStat = result[0]; // const latestStat = result[0];
if (result.length < 2) { // if (result.length < 2) {
const time = new Date(latestStat.updatedAt).getTime() - new Date(latestStat.createdAt).getTime(); // const time = new Date(latestStat.updatedAt).getTime() - new Date(latestStat.createdAt).getTime();
if (time < 1) { // // 1min
return 0; // if (time < 1000 * 60) {
} // return 0;
return (latestStat.shares * 4294967296) / (time / 1000); // }
} else { // return (parseFloat(latestStat.shares) * 4294967296) / (time / 1000);
const secondLatestStat = result[1]; // } else {
const time = new Date(latestStat.updatedAt).getTime() - new Date(secondLatestStat.createdAt).getTime(); // const secondLatestStat = result[1];
if (time < 1) { // const time = new Date(latestStat.updatedAt).getTime() - new Date(secondLatestStat.createdAt).getTime();
return 0; // // 1min
} // if (time < 1000 * 60) {
return ((latestStat.shares + secondLatestStat.shares) * 4294967296) / (time / 1000); // return 0;
} // }
// return ((parseFloat(latestStat.shares) + parseFloat(secondLatestStat.shares)) * 4294967296) / (time / 1000);
// }
} // }
public async getChartDataForSession(address: string, clientName: string, sessionId: string) { public async getChartDataForSession(clientId: string) {
var yesterday = new Date(new Date().getTime() - (24 * 60 * 60 * 1000)); var yesterday = new Date(new Date().getTime() - (24 * 60 * 60 * 1000));
const query = ` const query = `
SELECT SELECT
time label, time AS label,
(SUM(shares) * 4294967296) / 600 AS data (SUM(shares) * 4294967296) / 600 AS data
FROM FROM
client_statistics_entity AS entry client_statistics_entity AS entry
WHERE WHERE
entry.address = ? AND entry.clientName = ? AND entry.sessionId = ? AND entry.time > ${yesterday.getTime()} entry."clientId" = $1 AND entry.time > ${yesterday.getTime()}
GROUP BY GROUP BY
time time
ORDER BY ORDER BY
@@ -245,10 +249,10 @@ export class ClientStatisticsService {
LIMIT 144; LIMIT 144;
`; `;
const result = await this.clientStatisticsRepository.query(query, [address, clientName, sessionId]); const result = await this.clientStatisticsRepository.query(query, [clientId]);
return result.map(res => { return result.map(res => {
res.label = new Date(res.label).toISOString(); res.label = new Date(parseInt(res.label)).toISOString();
return res; return res;
}).slice(0, result.length - 1); }).slice(0, result.length - 1);
+23 -15
View File
@@ -1,25 +1,28 @@
import { Column, Entity, Index, PrimaryColumn } from 'typeorm'; import { Column, Entity, Index, OneToMany, PrimaryGeneratedColumn } from 'typeorm';
import { DateTimeTransformer } from '../utils/DateTimeTransformer'; import { ClientStatisticsEntity } from '../client-statistics/client-statistics.entity';
import { TrackedEntity } from '../utils/TrackedEntity.entity'; import { TrackedEntity } from '../utils/TrackedEntity.entity';
//https://www.sqlite.org/withoutrowid.html
//The WITHOUT ROWID optimization is likely to be helpful for tables that have non-integer @Entity()
// or composite (multi-column) PRIMARY KEYs and that do not store large strings or BLOBs. @Index("IDX_unique_nonce", { synchronize: false })
//WITHOUT ROWID tables work best when individual rows are not too large. @Index('idx_client_cleanup', ['id'], {
@Entity({ withoutRowid: true }) where: '"deletedAt" IS NULL',
@Index(['address', 'clientName', 'sessionId'], { unique: true }) // This is a partial (filtered) index — only indexes active clients
})
export class ClientEntity extends TrackedEntity { export class ClientEntity extends TrackedEntity {
@PrimaryGeneratedColumn('uuid')
id: string;
@PrimaryColumn({ length: 62, type: 'varchar' }) @Index()
@Column({ length: 62, type: 'varchar' })
address: string; address: string;
@PrimaryColumn({ length: 64, type: 'varchar' }) @Column({ length: 64, type: 'varchar' })
clientName: string; clientName: string;
@PrimaryColumn({ length: 8, type: 'varchar' }) @Column({ length: 8, type: 'varchar', })
sessionId: string; sessionId: string;
@@ -27,15 +30,20 @@ export class ClientEntity extends TrackedEntity {
userAgent: string; userAgent: string;
@Column({ type: 'timestamp' })
@Column({ type: 'datetime', transformer: new DateTimeTransformer() })
startTime: Date; startTime: Date;
@Column({ type: 'real', default: 0 }) @Column({ type: 'decimal', default: 0 })
bestDifficulty: number bestDifficulty: number
@Column({ default: 0 }) @Column({ default: 0, type: 'decimal' })
hashRate: number; hashRate: number;
@OneToMany(
() => ClientStatisticsEntity,
clientStatisticsEntity => clientStatisticsEntity.client
)
statistics: ClientStatisticsEntity[]
} }
+103 -56
View File
@@ -1,8 +1,6 @@
import { Injectable } from '@nestjs/common'; import { Injectable } from '@nestjs/common';
import { Interval } from '@nestjs/schedule'; import { InjectDataSource, InjectRepository } from '@nestjs/typeorm';
import { InjectRepository } from '@nestjs/typeorm'; import { DataSource, Repository } from 'typeorm';
import { BehaviorSubject, firstValueFrom } from 'rxjs';
import { ObjectLiteral, Repository } from 'typeorm';
import { ClientEntity } from './client.entity'; import { ClientEntity } from './client.entity';
@@ -11,70 +9,119 @@ import { ClientEntity } from './client.entity';
@Injectable() @Injectable()
export class ClientService { export class ClientService {
private heartbeatBulkUpdate: { [id: string]: { id: string, hashRate: number, updatedAt: Date } } = {};
public insertQueue: { result: BehaviorSubject<ObjectLiteral | null>, partialClient: Partial<ClientEntity> }[] = [];
constructor( constructor(
@InjectDataSource()
private dataSource: DataSource,
@InjectRepository(ClientEntity) @InjectRepository(ClientEntity)
private clientRepository: Repository<ClientEntity> private clientRepository: Repository<ClientEntity>
) { ) {
} }
@Interval(1000 * 5) // client.service.ts
public async insertClients() { public async killDeadClients(): Promise<boolean> {
const queueCopy = [...this.insertQueue]; const BATCH_SIZE = 10000;
this.insertQueue = [];
const results = await this.clientRepository.insert(queueCopy.map(c => c.partialClient)); return this.dataSource.transaction('READ COMMITTED', async (manager) => { // ← CHANGE HERE
const deadClients: { id: string }[] = await manager.query(`
SELECT id
FROM client_entity
WHERE "deletedAt" IS NULL
AND "updatedAt" < NOW() - INTERVAL '5 minutes'
ORDER BY id
LIMIT $1
FOR UPDATE SKIP LOCKED
`, [BATCH_SIZE]);
queueCopy.forEach((c, index) => { if (deadClients.length === 0) {
c.result.next(results.generatedMaps[index]); return false;
}
await manager
.createQueryBuilder()
.update(ClientEntity)
.set({ deletedAt: () => 'NOW()' })
.whereInIds(deadClients.map(c => c.id))
.execute();
console.log(`Killed ${deadClients.length} dead clients`);
return true;
}); });
} }
public async killDeadClients() { //public async heartbeat(id, hashRate: number, updatedAt: Date) {
var fiveMinutes = new Date(new Date().getTime() - (5 * 60 * 1000)).toISOString(); // return await this.clientRepository.update({ id }, { hashRate, deletedAt: null, updatedAt });
return await this.clientRepository
.createQueryBuilder()
.update(ClientEntity)
.set({ deletedAt: () => "DATETIME('now')" })
.where("deletedAt IS NULL AND updatedAt < DATETIME(:fiveMinutes)", { fiveMinutes })
.execute();
}
public async heartbeat(address: string, clientName: string, sessionId: string, hashRate: number, updatedAt: Date) {
return await this.clientRepository.update({ address, clientName, sessionId }, { hashRate, deletedAt: null, updatedAt });
}
// public async save(client: Partial<ClientEntity>) {
// return await this.clientRepository.save(client);
// } // }
public heartbeatBulkAsync(id, hashRate: number, updatedAt: Date) {
if (this.heartbeatBulkUpdate[id] != null) {
this.heartbeatBulkUpdate[id].hashRate = hashRate;
this.heartbeatBulkUpdate[id].updatedAt = updatedAt;
return;
}
this.heartbeatBulkUpdate[id] = { id, hashRate, updatedAt };
}
public async doBulkHeartbeatUpdate() {
if (Object.keys(this.heartbeatBulkUpdate).length < 1) {
console.log('No heartbeats to update.')
return;
}
const values = Object.entries(this.heartbeatBulkUpdate).map(([key, value]) => {
return `('${value.id}', ${value.hashRate}, NOW())`
}).join(',');
const query = `
DO $$
BEGIN
CREATE TEMP TABLE temp_heartbeats (
id UUID,
"hashRate" DECIMAL,
"updatedAt" TIMESTAMP
) ON COMMIT DROP;
INSERT INTO temp_heartbeats (id, "hashRate", "updatedAt")
VALUES ${values};
UPDATE "client_entity" ce
SET "hashRate" = th."hashRate",
"deletedAt" = NULL,
"updatedAt" = th."updatedAt"
FROM temp_heartbeats th
WHERE ce.id = th.id;
END;
$$;
`;
try {
await this.clientRepository.query(query);
console.log(`Bulk updated ${Object.keys(this.heartbeatBulkUpdate).length} client heartbeats`);
} catch (error) {
console.error('Bulk heartbeat failed:', error.message, 'Query:', query);
throw error;
}
this.heartbeatBulkUpdate = {};
}
public async insert(partialClient: Partial<ClientEntity>): Promise<ClientEntity> { public async insert(partialClient: Partial<ClientEntity>): Promise<ClientEntity> {
const insertResult = await this.clientRepository.insert(partialClient);
const result = new BehaviorSubject(null);
this.insertQueue.push({ result, partialClient });
// const insertResult = await this.clientRepository.insert(partialClient);
const generatedMap = await firstValueFrom(result);
const client = { const client = {
...partialClient, ...partialClient,
...generatedMap ...insertResult.generatedMaps[0]
}; };
return client as ClientEntity; return client as ClientEntity;
} }
public async delete(sessionId: string) { public async delete(id: string) {
return await this.clientRepository.softDelete({ sessionId }); return await this.clientRepository.softDelete({ id });
} }
public async deleteOldClients() { public async deleteOldClients() {
@@ -90,8 +137,8 @@ export class ClientService {
} }
public async updateBestDifficulty(sessionId: string, bestDifficulty: number) { public async updateBestDifficulty(id: string, bestDifficulty: number) {
return await this.clientRepository.update({ sessionId }, { bestDifficulty }); return await this.clientRepository.update({ id }, { bestDifficulty });
} }
public async connectedClientCount(): Promise<number> { public async connectedClientCount(): Promise<number> {
return await this.clientRepository.count(); return await this.clientRepository.count();
@@ -129,16 +176,16 @@ export class ClientService {
return await this.clientRepository.softDelete({}) return await this.clientRepository.softDelete({})
} }
public async getUserAgents() { // public async getUserAgents() {
const result = await this.clientRepository.createQueryBuilder('client') // const result = await this.clientRepository.createQueryBuilder('client')
.select('client.userAgent as userAgent') // .select('client.userAgent as "userAgent"')
.addSelect('COUNT(client.userAgent)', 'count') // .addSelect('COUNT(client.userAgent)', 'count')
.addSelect('MAX(client.bestDifficulty)', 'bestDifficulty') // .addSelect('MAX(client.bestDifficulty)', 'bestDifficulty')
.addSelect('SUM(client.hashRate)', 'totalHashRate') // .addSelect('SUM(client.hashRate)', 'totalHashRate')
.groupBy('client.userAgent') // .groupBy('client.userAgent')
.orderBy('count', 'DESC') // .orderBy('"totalHashRate"', 'DESC')
.getRawMany(); // .getRawMany();
return result; // return result;
} // }
} }
+14
View File
@@ -0,0 +1,14 @@
import { Column, Entity, PrimaryGeneratedColumn } from 'typeorm';
@Entity()
export class HomeGraphEntity {
@PrimaryGeneratedColumn({type: 'bigint'})
id: number;
@Column({ type: 'bigint' })
label: number;
@Column({ type: 'bigint' })
data: number;
}
+15
View File
@@ -0,0 +1,15 @@
import { Global, Module } from '@nestjs/common';
import { TypeOrmModule } from '@nestjs/typeorm';
import { HomeGraphEntity } from './home-graph.entity';
import { HomeGraphService } from './home-graph.service';
@Global()
@Module({
imports: [TypeOrmModule.forFeature([HomeGraphEntity])],
providers: [HomeGraphService],
exports: [TypeOrmModule, HomeGraphService],
})
export class HomeGraphModule { }
+45
View File
@@ -0,0 +1,45 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { HomeGraphEntity } from './home-graph.entity';
@Injectable()
export class HomeGraphService {
constructor(
@InjectRepository(HomeGraphEntity) private homeGraphRepository: Repository<HomeGraphEntity>
) {
}
public async getLatestTime(): Promise<Date> {
const result = await this.homeGraphRepository.createQueryBuilder('entity')
.select('MAX(entity.label)', 'maxNumber')
.getRawOne();
if (result.maxNumber == null) {
return new Date(new Date().getTime() - (24 * 60 * 60 * 1000));
}
return new Date(parseInt(result.maxNumber));
}
public async save(graph: { label: number, data: number }[]) {
await this.homeGraphRepository.insert(graph);
}
public async getChartDataForSite(limit: number = 144 * 7) {
const records = await this.homeGraphRepository
.createQueryBuilder('homeGraph')
.orderBy('homeGraph.id', 'DESC')
.limit(limit)
.getRawMany();
return records.map(res => {
return {
label: new Date(parseInt(res.homeGraph_label)).toISOString(),
data: res.homeGraph_data
}
});
}
}
-3
View File
@@ -6,9 +6,6 @@ export class RpcBlockEntity {
@PrimaryColumn() @PrimaryColumn()
blockHeight: number; blockHeight: number;
@Column({ nullable: true })
lockedBy?: string;
@Column({ nullable: true }) @Column({ nullable: true })
data?: string; data?: string;
} }
+17 -6
View File
@@ -1,6 +1,6 @@
import { Injectable } from '@nestjs/common'; import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm'; import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm'; import { InsertResult, Repository } from 'typeorm';
import { RpcBlockEntity } from './rpc-block.entity'; import { RpcBlockEntity } from './rpc-block.entity';
@@ -12,17 +12,28 @@ export class RpcBlockService {
) { ) {
} }
public getBlock(blockHeight: number) { public getSavedBlockTemplate(blockHeight: number) {
return this.rpcBlockRepository.findOne({ return this.rpcBlockRepository.findOne({
where: { blockHeight } where: { blockHeight }
}); });
} }
public lockBlock(blockHeight: number, process: string) { public saveBlock(blockHeight: number, data: string): Promise<InsertResult> {
return this.rpcBlockRepository.save({ blockHeight, data: null, lockedBy: process }); return this.rpcBlockRepository.upsert({ blockHeight, data }, ['blockHeight']);
} }
public saveBlock(blockHeight: number, data: string) { public async deleteOldBlocks() {
return this.rpcBlockRepository.update(blockHeight, { data }) const result = await this.rpcBlockRepository.createQueryBuilder('entity')
.select('MAX(entity.blockHeight)', 'maxNumber')
.getRawOne();
const newestBlock = result ? result.maxNumber : null;
await this.rpcBlockRepository.createQueryBuilder()
.delete()
.where('"blockHeight" < :newestBlock', { newestBlock })
.execute();
return;
} }
} }
@@ -5,7 +5,7 @@ import { TrackedEntity } from '../utils/TrackedEntity.entity';
@Entity() @Entity()
export class TelegramSubscriptionsEntity extends TrackedEntity { export class TelegramSubscriptionsEntity extends TrackedEntity {
@PrimaryGeneratedColumn() @PrimaryGeneratedColumn({type: 'bigint'})
id: number; id: number;
@Index() @Index()
+3 -5
View File
@@ -1,14 +1,12 @@
import { CreateDateColumn, DeleteDateColumn, UpdateDateColumn } from 'typeorm'; import { CreateDateColumn, DeleteDateColumn, UpdateDateColumn } from 'typeorm';
import { DateTimeTransformer } from './DateTimeTransformer';
export abstract class TrackedEntity { export abstract class TrackedEntity {
@DeleteDateColumn({ nullable: true, type: 'datetime', transformer: new DateTimeTransformer() }) @DeleteDateColumn({ nullable: true, type: 'timestamp' })
public deletedAt?: Date; public deletedAt?: Date;
@CreateDateColumn({ type: 'datetime', transformer: new DateTimeTransformer() }) @CreateDateColumn({ type: 'timestamp' })
public createdAt?: Date public createdAt?: Date
@UpdateDateColumn({ type: 'datetime', transformer: new DateTimeTransformer() }) @UpdateDateColumn({ type: 'timestamp' })
public updatedAt?: Date public updatedAt?: Date
} }
+43 -11
View File
@@ -3,10 +3,14 @@ import { Controller, Get, Inject } from '@nestjs/common';
import { Cache } from 'cache-manager'; import { Cache } from 'cache-manager';
import { firstValueFrom } from 'rxjs'; import { firstValueFrom } from 'rxjs';
import { UserAgentReportService } from './ORM/_views/user-agent-report/user-agent-report.service';
import { AddressSettingsService } from './ORM/address-settings/address-settings.service';
import { BlocksService } from './ORM/blocks/blocks.service'; import { BlocksService } from './ORM/blocks/blocks.service';
import { ClientStatisticsService } from './ORM/client-statistics/client-statistics.service'; import { ClientStatisticsService } from './ORM/client-statistics/client-statistics.service';
import { ClientService } from './ORM/client/client.service'; import { ClientService } from './ORM/client/client.service';
import { HomeGraphService } from './ORM/home-graph/home-graph.service';
import { BitcoinRpcService } from './services/bitcoin-rpc.service'; import { BitcoinRpcService } from './services/bitcoin-rpc.service';
import { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-report.view';
@Controller() @Controller()
export class AppController { export class AppController {
@@ -18,7 +22,10 @@ export class AppController {
private readonly clientService: ClientService, private readonly clientService: ClientService,
private readonly clientStatisticsService: ClientStatisticsService, private readonly clientStatisticsService: ClientStatisticsService,
private readonly blocksService: BlocksService, private readonly blocksService: BlocksService,
private readonly bitcoinRpcService: BitcoinRpcService private readonly bitcoinRpcService: BitcoinRpcService,
private readonly homeGraphService: HomeGraphService,
private readonly addressSettingsService: AddressSettingsService,
private readonly userAgentReportService: UserAgentReportService
) { } ) { }
@Get('info') @Get('info')
@@ -34,16 +41,42 @@ export class AppController {
const blockData = await this.blocksService.getFoundBlocks(); const blockData = await this.blocksService.getFoundBlocks();
const userAgents = await this.clientService.getUserAgents(); const highScores = await this.addressSettingsService.getHighScores();
const other: {
count: number,
bestDifficulty: number,
totalHashRate: number;
} = {
count: 0,
bestDifficulty: 0,
totalHashRate: 0
};
const userAgents: UserAgentReportView[] = (await this.userAgentReportService.getReport()).reduce((pre, cur, idx, arr) => {
// If less than 10Th/s and less than 100 devices, add to 'other'
if (parseInt(cur.totalHashRate) < 10000000000000 && parseInt(cur.count) < 200) {
other.totalHashRate += parseFloat(cur.totalHashRate);
other.count += parseInt(cur.count);
if (other.bestDifficulty < cur.bestDifficulty) {
other.bestDifficulty = cur.bestDifficulty;
}
} else {
pre.push(cur);
}
return pre;
}, []);
userAgents.push({ userAgent: 'Other', count: other.count.toString(), bestDifficulty: other.bestDifficulty, totalHashRate: other.totalHashRate.toString() })
const data = { const data = {
blockData, blockData,
userAgents, userAgents,
highScores,
uptime: this.uptime uptime: this.uptime
}; };
//1 min //5 min
await this.cacheManager.set(CACHE_KEY, data, 1 * 60 * 1000); await this.cacheManager.set(CACHE_KEY, data, 5 * 60 * 1000);
return data; return data;
@@ -59,13 +92,13 @@ export class AppController {
return cachedResult; return cachedResult;
} }
const userAgents = await this.userAgentReportService.getReport();
const userAgents = await this.clientService.getUserAgents();
const totalHashRate = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.totalHashRate), 0); const totalHashRate = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.totalHashRate), 0);
const totalMiners = userAgents.reduce((acc, userAgent) => acc + parseInt(userAgent.count), 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 blocksFound = await this.blocksService.getFoundBlocks();
const data = { const data = {
totalHashRate, totalHashRate,
blockHeight, blockHeight,
@@ -82,8 +115,7 @@ export class AppController {
@Get('network') @Get('network')
public async network() { public async network() {
const miningInfo = await firstValueFrom(this.bitcoinRpcService.newBlock$); return this.bitcoinRpcService.miningInfo;
return miningInfo;
} }
@Get('info/chart') @Get('info/chart')
@@ -97,7 +129,7 @@ export class AppController {
return cachedResult; return cachedResult;
} }
const chartData = await this.clientStatisticsService.getChartDataForSite(); const chartData = await this.homeGraphService.getChartDataForSite();
//10 min //10 min
await this.cacheManager.set(CACHE_KEY, chartData, 10 * 60 * 1000); await this.cacheManager.set(CACHE_KEY, chartData, 10 * 60 * 1000);
+44 -13
View File
@@ -1,7 +1,7 @@
import { HttpModule } from '@nestjs/axios'; import { HttpModule } from '@nestjs/axios';
import { CacheModule } from '@nestjs/cache-manager'; import { CacheModule } from '@nestjs/cache-manager';
import { Module } from '@nestjs/common'; import { Module } from '@nestjs/common';
import { ConfigModule } from '@nestjs/config'; import { ConfigModule, ConfigService } from '@nestjs/config';
import { ScheduleModule } from '@nestjs/schedule'; import { ScheduleModule } from '@nestjs/schedule';
import { TypeOrmModule } from '@nestjs/typeorm'; import { TypeOrmModule } from '@nestjs/typeorm';
@@ -9,11 +9,22 @@ import { AppController } from './app.controller';
import { AddressController } from './controllers/address/address.controller'; import { AddressController } from './controllers/address/address.controller';
import { ClientController } from './controllers/client/client.controller'; import { ClientController } from './controllers/client/client.controller';
import { BitcoinAddressValidator } from './models/validators/bitcoin-address.validator'; import { BitcoinAddressValidator } from './models/validators/bitcoin-address.validator';
import { UniqueNonceIndex } from './ORM/_migrations/UniqueNonceIndex';
import { UserAgentReportModule } from './ORM/_views/user-agent-report/user-agent-report.module';
import { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-report.view';
import { AddressSettingsEntity } from './ORM/address-settings/address-settings.entity';
import { AddressSettingsModule } from './ORM/address-settings/address-settings.module'; import { AddressSettingsModule } from './ORM/address-settings/address-settings.module';
import { BlocksEntity } from './ORM/blocks/blocks.entity';
import { BlocksModule } from './ORM/blocks/blocks.module'; import { BlocksModule } from './ORM/blocks/blocks.module';
import { ClientStatisticsEntity } from './ORM/client-statistics/client-statistics.entity';
import { ClientStatisticsModule } from './ORM/client-statistics/client-statistics.module'; import { ClientStatisticsModule } from './ORM/client-statistics/client-statistics.module';
import { ClientEntity } from './ORM/client/client.entity';
import { ClientModule } from './ORM/client/client.module'; import { ClientModule } from './ORM/client/client.module';
import { HomeGraphEntity } from './ORM/home-graph/home-graph.entity';
import { HomeGraphModule } from './ORM/home-graph/home-graph.module';
import { RpcBlockEntity } from './ORM/rpc-block/rpc-block.entity';
import { RpcBlocksModule } from './ORM/rpc-block/rpc-block.module'; import { RpcBlocksModule } from './ORM/rpc-block/rpc-block.module';
import { TelegramSubscriptionsEntity } from './ORM/telegram-subscriptions/telegram-subscriptions.entity';
import { TelegramSubscriptionsModule } from './ORM/telegram-subscriptions/telegram-subscriptions.module'; import { TelegramSubscriptionsModule } from './ORM/telegram-subscriptions/telegram-subscriptions.module';
import { AppService } from './services/app.service'; import { AppService } from './services/app.service';
import { BitcoinRpcService } from './services/bitcoin-rpc.service'; import { BitcoinRpcService } from './services/bitcoin-rpc.service';
@@ -21,7 +32,6 @@ import { BraiinsService } from './services/braiins.service';
import { BTCPayService } from './services/btc-pay.service'; import { BTCPayService } from './services/btc-pay.service';
import { DiscordService } from './services/discord.service'; import { DiscordService } from './services/discord.service';
import { NotificationService } from './services/notification.service'; import { NotificationService } from './services/notification.service';
import { ProxyService } from './services/proxy.service';
import { StratumV1JobsService } from './services/stratum-v1-jobs.service'; import { StratumV1JobsService } from './services/stratum-v1-jobs.service';
import { StratumV1Service } from './services/stratum-v1.service'; import { StratumV1Service } from './services/stratum-v1.service';
import { TelegramService } from './services/telegram.service'; import { TelegramService } from './services/telegram.service';
@@ -33,21 +43,43 @@ const ORMModules = [
AddressSettingsModule, AddressSettingsModule,
TelegramSubscriptionsModule, TelegramSubscriptionsModule,
BlocksModule, BlocksModule,
RpcBlocksModule RpcBlocksModule,
HomeGraphModule,
UserAgentReportModule
] ]
@Module({ @Module({
imports: [ imports: [
ConfigModule.forRoot(), ConfigModule.forRoot(),
TypeOrmModule.forRoot({ TypeOrmModule.forRootAsync({
type: 'sqlite', useFactory: (configService: ConfigService) => {
database: './DB/public-pool.sqlite', return {
synchronize: true, type: 'postgres',
autoLoadEntities: true, host: configService.get('DB_HOST'),
logging: false, port: parseInt(configService.get('DB_PORT')),
enableWAL: true, username: configService.get('DB_USERNAME'),
busyTimeout: 30 * 1000, password: configService.get('DB_PASSWORD'),
database: configService.get('DB_DATABASE'),
entities: [
ClientEntity,
AddressSettingsEntity,
BlocksEntity,
ClientStatisticsEntity,
RpcBlockEntity,
TelegramSubscriptionsEntity,
HomeGraphEntity,
UserAgentReportView
],
synchronize: configService.get('PRODUCTION') != 'true',
logging: false,
poolSize: 10,
migrations: [
UniqueNonceIndex
]
}
},
imports: [ConfigModule],
inject: [ConfigService]
}), }),
CacheModule.register(), CacheModule.register(),
ScheduleModule.forRoot(), ScheduleModule.forRoot(),
@@ -68,7 +100,6 @@ const ORMModules = [
NotificationService, NotificationService,
BitcoinAddressValidator, BitcoinAddressValidator,
StratumV1JobsService, StratumV1JobsService,
ProxyService,
BTCPayService, BTCPayService,
BraiinsService BraiinsService
], ],
+2 -2
View File
@@ -30,7 +30,7 @@ export class ClientController {
return { return {
sessionId: worker.sessionId, sessionId: worker.sessionId,
name: worker.clientName, name: worker.clientName,
bestDifficulty: worker.bestDifficulty.toFixed(2), bestDifficulty: parseFloat(worker.bestDifficulty as any).toFixed(2),
hashRate: worker.hashRate, hashRate: worker.hashRate,
startTime: worker.startTime, startTime: worker.startTime,
lastSeen: worker.updatedAt lastSeen: worker.updatedAt
@@ -75,7 +75,7 @@ export class ClientController {
if (worker == null) { if (worker == null) {
return new NotFoundException(); return new NotFoundException();
} }
const chartData = await this.clientStatisticsService.getChartDataForSession(worker.address, worker.clientName, worker.sessionId); const chartData = await this.clientStatisticsService.getChartDataForSession(worker.id);
return { return {
sessionId: worker.sessionId, sessionId: worker.sessionId,
+54 -11
View File
@@ -3,38 +3,41 @@ import { NestFactory } from '@nestjs/core';
import { FastifyAdapter, NestFastifyApplication } from '@nestjs/platform-fastify'; import { FastifyAdapter, NestFastifyApplication } from '@nestjs/platform-fastify';
import * as bitcoinjs from 'bitcoinjs-lib'; import * as bitcoinjs from 'bitcoinjs-lib';
import { useContainer } from 'class-validator'; import { useContainer } from 'class-validator';
import { readFileSync } from 'fs'; import { readFileSync, watch } from 'fs';
import * as path from 'path';
import * as ecc from 'tiny-secp256k1'; import * as ecc from 'tiny-secp256k1';
import { AppModule } from './app.module'; import { AppModule } from './app.module';
async function bootstrap() { async function bootstrap() {
if (process.env.API_PORT == null) { if (process.env.API_PORT == null) {
console.error('It appears your environment is not configured, create and populate an .env file.'); console.error('It appears your environment is not configured, create and populate an .env file.');
return; return;
} }
let options = {}; const secure = process.env.API_SECURE?.toLowerCase() === 'true';
const secure = process.env.API_SECURE?.toLowerCase() == 'true'; const currentDirectory = process.cwd();
const keyPath = path.join(currentDirectory, 'secrets', 'key.pem');
const certPath = path.join(currentDirectory, 'secrets', 'cert.pem');
let options: any = {};
if (secure) { if (secure) {
const currentDirectory = process.cwd();
options = { options = {
https: { https: {
key: readFileSync(`${currentDirectory}/secrets/key.pem`), key: readFileSync(keyPath),
cert: readFileSync(`${currentDirectory}/secrets/cert.pem`), cert: readFileSync(certPath),
} }
}; };
} }
const app = await NestFactory.create<NestFastifyApplication>(AppModule, new FastifyAdapter(options)); const app = await NestFactory.create<NestFastifyApplication>(AppModule, new FastifyAdapter(options));
app.setGlobalPrefix('api') app.setGlobalPrefix('api');
app.useGlobalPipes( app.useGlobalPipes(
new ValidationPipe({ new ValidationPipe({
transform: true, transform: true,
whitelist: true, whitelist: true,
forbidNonWhitelisted: true, //forbidNonWhitelisted: true,
forbidUnknownValues: true //forbidUnknownValues: true
}), }),
); );
@@ -51,13 +54,53 @@ async function bootstrap() {
app.enableCors(); app.enableCors();
useContainer(app.select(AppModule), { fallbackOnErrors: true }); useContainer(app.select(AppModule), { fallbackOnErrors: true });
//Taproot // Taproot
bitcoinjs.initEccLib(ecc); bitcoinjs.initEccLib(ecc);
await app.listen(process.env.API_PORT, '0.0.0.0', (err, address) => { await app.listen(process.env.API_PORT, '0.0.0.0', (err, address) => {
if (err) {
console.error(err);
process.exit(1);
}
console.log(`API listening on ${address}`); console.log(`API listening on ${address}`);
}); });
// --- Live-reload TLS certs/keys when they change on disk ---
if (secure) {
// Fastify's underlying Node https server
const server: any = app.getHttpServer();
// Guard: only HTTPS servers expose setSecureContext
if (typeof server?.setSecureContext === 'function') {
let reloadTimer: NodeJS.Timeout | null = null;
const scheduleReload = () => {
if (reloadTimer) clearTimeout(reloadTimer);
// Debounce multiple fs events during a single write/replace
reloadTimer = setTimeout(() => {
try {
const key = readFileSync(keyPath);
const cert = readFileSync(certPath);
server.setSecureContext({ key, cert });
console.log(`[TLS] Reloaded certificate @ ${new Date().toISOString()}`);
} catch (e) {
console.error('[TLS] Failed to reload certificate:', e);
}
}, 500);
};
// Watch both files; handle 'change' and 'rename' (rename often fired on atomic replace)
try {
watch(keyPath, { persistent: true }, scheduleReload);
watch(certPath, { persistent: true }, scheduleReload);
console.log('[TLS] Watching cert/key for changes');
} catch (e) {
console.error('[TLS] Failed to watch cert/key files:', e);
}
} else {
console.warn('[TLS] Dynamic cert reload not available (non-HTTPS server?)');
}
}
} }
bootstrap(); bootstrap();
+3 -2
View File
@@ -18,7 +18,7 @@ export class MiningJob {
public jobTemplateId: string; public jobTemplateId: string;
public networkDifficulty: number; public networkDifficulty: number;
public creation: number;
constructor( constructor(
private network: bitcoinjs.networks.Network, private network: bitcoinjs.networks.Network,
@@ -27,6 +27,7 @@ export class MiningJob {
jobTemplate: IJobTemplate jobTemplate: IJobTemplate
) { ) {
this.creation = new Date().getTime();
this.jobTemplateId = jobTemplate.blockData.id; this.jobTemplateId = jobTemplate.blockData.id;
this.coinbaseTransaction = this.createCoinbaseTransaction(payoutInformation, jobTemplate.blockData.coinbasevalue); this.coinbaseTransaction = this.createCoinbaseTransaction(payoutInformation, jobTemplate.blockData.coinbasevalue);
@@ -39,7 +40,7 @@ export class MiningJob {
// 32-byte - Commitment hash: Double-SHA256(witness root hash|witness reserved value) // 32-byte - Commitment hash: Double-SHA256(witness root hash|witness reserved value)
// 39th byte onwards: Optional data with no consensus meaning // 39th byte onwards: Optional data with no consensus meaning
const extra = Buffer.from('\\public-pool\\'); const extra = Buffer.from('Public-Pool');
// Encode the block height // Encode the block height
// https://github.com/bitcoin/bips/blob/master/bip-0034.mediawiki // https://github.com/bitcoin/bips/blob/master/bip-0034.mediawiki
-235
View File
@@ -1,235 +0,0 @@
import Big from 'big.js';
import { Block } from 'bitcoinjs-lib';
import * as bitcoinjs from 'bitcoinjs-lib';
import { plainToInstance } from 'class-transformer';
import { validate, ValidatorOptions } from 'class-validator';
import { Socket } from 'net';
import { eRequestMethod } from './enums/eRequestMethod';
import { eResponseMethod } from './enums/eResponseMethod';
import { IMiningNotify } from './stratum-messages/IMiningNotify';
import { MiningSubmitMessage } from './stratum-messages/MiningSubmitMessage';
export class ProxyClient {
public braiinsSocket: Socket;
public serverDifficulty: number;
public jobs: { [jobId: string]: IMiningNotify } = {};
private extraNonce: string;
constructor(public socket: Socket) {
this.braiinsSocket = new Socket();
this.braiinsSocket.connect(3333, 'us-east.stratum.braiins.com', () => {
console.log('Connected to braiins');
});
this.braiinsSocket.on('data', (data) => {
data.toString()
.split('\n')
.filter(m => m.length > 0)
.forEach(async (m) => {
try {
await this.toClient(m);
} catch (e) {
await socket.end();
console.error(e);
}
})
});
this.braiinsSocket.on('close', () => {
console.log('Braiins closed connection');
socket.end();
});
socket.on('error', async (error: Error) => { });
socket.on('data', (data: Buffer) => {
data.toString()
.split('\n')
.filter(m => m.length > 0)
.forEach(async (m) => {
try {
await this.toServer(m);
} catch (e) {
await socket.end();
console.error(e);
}
})
});
socket.on('close', () => {
console.log('Client closed connection');
this.braiinsSocket.end();
})
}
private async toClient(message: string) {
let parsedMessage: IMiningNotify | any;
try {
parsedMessage = JSON.parse(message);
} catch (e) {
//console.log("Invalid JSON");
await this.socket.end();
return;
}
if (parsedMessage.method == null && Array.isArray(parsedMessage.result)) {
this.extraNonce = parsedMessage.result[1];
}
switch (parsedMessage.method) {
case eResponseMethod.SET_DIFFICULTY: {
this.serverDifficulty = parsedMessage.params[0];
parsedMessage.params[0] = 512;
message = JSON.stringify(parsedMessage);
break;
}
case eResponseMethod.MINING_NOTIFY: {
// clear jobs
if ((parsedMessage as IMiningNotify).params[8]) {
this.jobs = {};
}
this.jobs[parsedMessage.params[0]] = parsedMessage;
break;
}
}
console.log('Server:');
console.log(message);
this.socket.write(message + `\n`);
}
private async toServer(message: string) {
let parsedMessage = null;
try {
parsedMessage = JSON.parse(message);
} catch (e) {
//console.log("Invalid JSON");
await this.socket.end();
return;
}
switch (parsedMessage.method) {
case eRequestMethod.AUTHORIZE: {
parsedMessage.params[0] = 'battlechicken';
message = JSON.stringify(parsedMessage);
break;
}
case eRequestMethod.SUBMIT: {
parsedMessage.params[0] = 'battlechicken';
message = JSON.stringify(parsedMessage);
const miningSubmitMessage = plainToInstance(
MiningSubmitMessage,
parsedMessage,
);
const validatorOptions: ValidatorOptions = {
whitelist: true,
forbidNonWhitelisted: true,
};
const errors = await validate(miningSubmitMessage, validatorOptions);
if (errors.length === 0) {
this.checkSubmission(miningSubmitMessage);
}
else {
console.log('parsing error')
}
break;
}
}
console.log('Client:');
console.log(message);
this.braiinsSocket.write(message + '\n');
}
private checkSubmission(submission: MiningSubmitMessage) {
console.log('Submission:');
console.log(submission);
const job = this.jobs[submission.jobId];
if (job == null) {
console.log('Job not found');
return;
}
const versionMask = parseInt(submission.versionMask, 16);
const block = new Block();
block.version = parseInt(job.params[5], 16);
if (submission.versionMask !== undefined && versionMask != 0) {
block.version = (block.version ^ versionMask);
}
const prevHash = this.swapEndianWords(job.params[1]);
block.prevHash = Buffer.from(prevHash, 'hex');
const coinbase = Buffer.from(`${job.params[2]}${this.extraNonce}${submission.extraNonce2}${job.params[3]}`, 'hex');
const coinbaseHash = bitcoinjs.crypto.hash256(coinbase);
block.merkleRoot = this.calculateMerkleRootHash(coinbaseHash, job.params[4]);
block.timestamp = parseInt(submission.ntime, 16);
block.bits = parseInt(job.params[6], 16);
block.nonce = parseInt(submission.nonce, 16);
const header = block.toBuffer(true);
const diff = this.calculateDifficulty(header);
console.log(`DIFFICULTY: ${diff.submissionDifficulty}`)
}
private calculateMerkleRootHash(newRoot: Buffer, merkleBranches: string[]): Buffer {
const bothMerkles = Buffer.alloc(64);
bothMerkles.set(newRoot);
for (let i = 0; i < merkleBranches.length; i++) {
bothMerkles.set(Buffer.from(merkleBranches[i], 'hex'), 32);
newRoot = bitcoinjs.crypto.hash256(bothMerkles);
bothMerkles.set(newRoot);
}
return bothMerkles.subarray(0, 32)
}
private calculateDifficulty(header: Buffer): { submissionDifficulty: number, submissionHash: string } {
const hashResult = bitcoinjs.crypto.hash256(header);
let s64 = this.le256todouble(hashResult);
const truediffone = Big('26959535291011309493156476344723991336010898738574164086137773096960');
const difficulty = truediffone.div(s64.toString());
return { submissionDifficulty: difficulty.toNumber(), submissionHash: hashResult.toString('hex') };
}
private le256todouble(target: Buffer): bigint {
const number = target.reduceRight((acc, byte) => {
// Shift the number 8 bits to the left and OR with the current byte
return (acc << BigInt(8)) | BigInt(byte);
}, BigInt(0));
return number;
}
private swapEndianWords(str: string) {
const hexGroups = str.match(/.{1,8}/g);
// Reverse each group and concatenate them
const reversedHexString = hexGroups.reduce((pre, cur, indx, arr) => {
const reversed = cur.match(/.{2}/g).reverse();
return `${pre}${reversed.join('')}`;
}, '');
return reversedHexString;
}
}
+1 -1
View File
@@ -202,7 +202,7 @@ describe('StratumV1Client', () => {
expect(socket.on).toHaveBeenCalled(); expect(socket.on).toHaveBeenCalled();
socketEmitter(Buffer.from(MockRecording1.MINING_SUGGEST_DIFFICULTY)); socketEmitter(Buffer.from(MockRecording1.MINING_SUGGEST_DIFFICULTY));
await new Promise((r) => setTimeout(r, 1)); await new Promise((r) => setTimeout(r, 1));
expect(socket.write).toHaveBeenCalledWith(`{"id":4,"method":"mining.set_difficulty","params":[512]}\n`, expect.any(Function)); expect(socket.write).toHaveBeenCalledWith(`{"id":null,"method":"mining.set_difficulty","params":[512]}\n`, expect.any(Function));
}); });
it('should set difficulty', async () => { it('should set difficulty', async () => {
+126 -70
View File
@@ -31,7 +31,7 @@ import { StratumV1ClientStatistics } from './StratumV1ClientStatistics';
export class StratumV1Client { export class StratumV1Client {
private clientSubscription: SubscriptionMessage; public clientSubscription: SubscriptionMessage;
private clientConfiguration: ConfigurationMessage; private clientConfiguration: ConfigurationMessage;
private clientAuthorization: AuthorizationMessage; private clientAuthorization: AuthorizationMessage;
private clientSuggestedDifficulty: SuggestDifficulty; private clientSuggestedDifficulty: SuggestDifficulty;
@@ -41,16 +41,19 @@ export class StratumV1Client {
private statistics: StratumV1ClientStatistics; private statistics: StratumV1ClientStatistics;
private stratumInitialized = false; private stratumInitialized = false;
private usedSuggestedDifficulty = false; private usedSuggestedDifficulty = false;
private sessionDifficulty: number = 16384; private sessionDifficulty: number = 100000;
private entity: ClientEntity; private clientEntity: ClientEntity;
private creatingEntity: Promise<void>; private creatingEntity: Promise<void>;
public extraNonceAndSessionId: string; public extraNonceAndSessionId: string;
public sessionStart: Date; public sessionStart: Date;
public noFee: boolean; //public noFee: boolean;
public hashRate: number; //public hashRate: number = 0;
private buffer: string = '';
private miningSubmissionHashes = new Set<string>()
constructor( constructor(
public readonly socket: Socket, public readonly socket: Socket,
@@ -65,8 +68,11 @@ export class StratumV1Client {
) { ) {
this.socket.on('data', (data: Buffer) => { this.socket.on('data', (data: Buffer) => {
data.toString() this.buffer += data.toString();
.split('\n') let lines = this.buffer.split('\n');
this.buffer = lines.pop() || ''; // Save the last part of the data (incomplete line) to the buffer
lines
.filter(m => m.length > 0) .filter(m => m.length > 0)
.forEach(async (m) => { .forEach(async (m) => {
try { try {
@@ -75,7 +81,7 @@ export class StratumV1Client {
await this.socket.end(); await this.socket.end();
console.error(e); console.error(e);
} }
}) });
}); });
@@ -83,7 +89,9 @@ export class StratumV1Client {
public async destroy() { public async destroy() {
await this.clientService.delete(this.extraNonceAndSessionId); if (this.clientEntity?.id) {
await this.clientService.delete(this.clientEntity.id);
}
if (this.stratumSubscription != null) { if (this.stratumSubscription != null) {
this.stratumSubscription.unsubscribe(); this.stratumSubscription.unsubscribe();
@@ -126,7 +134,7 @@ export class StratumV1Client {
const validatorOptions: ValidatorOptions = { const validatorOptions: ValidatorOptions = {
whitelist: true, whitelist: true,
forbidNonWhitelisted: true, //forbidNonWhitelisted: true,
}; };
const errors = await validate(subscriptionMessage, validatorOptions); const errors = await validate(subscriptionMessage, validatorOptions);
@@ -137,7 +145,7 @@ export class StratumV1Client {
this.sessionStart = new Date(); this.sessionStart = new Date();
this.statistics = new StratumV1ClientStatistics(this.clientStatisticsService); this.statistics = new StratumV1ClientStatistics(this.clientStatisticsService);
this.extraNonceAndSessionId = this.getRandomHexString(); this.extraNonceAndSessionId = this.getRandomHexString();
console.log(`New client ID: : ${this.extraNonceAndSessionId}, ${this.socket.remoteAddress}:${this.socket.remotePort}`); //console.log(`New client ID: : ${this.extraNonceAndSessionId}, ${this.socket.remoteAddress}:${this.socket.remotePort}`);
} }
this.clientSubscription = subscriptionMessage; this.clientSubscription = subscriptionMessage;
@@ -170,7 +178,7 @@ export class StratumV1Client {
const validatorOptions: ValidatorOptions = { const validatorOptions: ValidatorOptions = {
whitelist: true, whitelist: true,
forbidNonWhitelisted: true, //forbidNonWhitelisted: true,
}; };
const errors = await validate(configurationMessage, validatorOptions); const errors = await validate(configurationMessage, validatorOptions);
@@ -184,7 +192,7 @@ export class StratumV1Client {
} }
} else { } else {
console.error('Configuration validation error'); console.log('Configuration validation error');
const err = new StratumErrorMessage( const err = new StratumErrorMessage(
configurationMessage.id, configurationMessage.id,
eStratumErrorCode.OtherUnknown, eStratumErrorCode.OtherUnknown,
@@ -208,25 +216,27 @@ export class StratumV1Client {
const validatorOptions: ValidatorOptions = { const validatorOptions: ValidatorOptions = {
whitelist: true, whitelist: true,
forbidNonWhitelisted: true, //forbidNonWhitelisted: true,
}; };
const errors = await validate(authorizationMessage, validatorOptions); const errors = await validate(authorizationMessage, validatorOptions);
if (errors.length === 0) { if (errors.length === 0) {
this.clientAuthorization = authorizationMessage; this.clientAuthorization = authorizationMessage;
if (this.clientSuggestedDifficulty == null && this.clientAuthorization.startingDiff != null && this.clientAuthorization.startingDiff > this.sessionDifficulty) {
this.sessionDifficulty = this.clientAuthorization.startingDiff;
}
const success = await this.write(JSON.stringify(this.clientAuthorization.response()) + '\n'); const success = await this.write(JSON.stringify(this.clientAuthorization.response()) + '\n');
if (!success) { if (!success) {
return; return;
} }
} else { } else {
console.error('Authorization validation error');
const err = new StratumErrorMessage( const err = new StratumErrorMessage(
authorizationMessage.id, authorizationMessage.id,
eStratumErrorCode.OtherUnknown, eStratumErrorCode.OtherUnknown,
'Authorization validation error', 'Authorization validation error',
errors).response(); errors).response();
console.error(err); //console.log(err);
const success = await this.write(err); const success = await this.write(err);
if (!success) { if (!success) {
return; return;
@@ -247,7 +257,7 @@ export class StratumV1Client {
const validatorOptions: ValidatorOptions = { const validatorOptions: ValidatorOptions = {
whitelist: true, whitelist: true,
forbidNonWhitelisted: true, //forbidNonWhitelisted: true,
}; };
const errors = await validate(suggestDifficultyMessage, validatorOptions); const errors = await validate(suggestDifficultyMessage, validatorOptions);
@@ -279,7 +289,7 @@ export class StratumV1Client {
case eRequestMethod.SUBMIT: { case eRequestMethod.SUBMIT: {
if (this.stratumInitialized == false) { if (this.stratumInitialized == false) {
console.log('Submit before initalized'); //console.log('Submit before initalized');
await this.socket.end(); await this.socket.end();
return; return;
} }
@@ -292,7 +302,7 @@ export class StratumV1Client {
const validatorOptions: ValidatorOptions = { const validatorOptions: ValidatorOptions = {
whitelist: true, whitelist: true,
forbidNonWhitelisted: true, //forbidNonWhitelisted: true,
}; };
const errors = await validate(miningSubmitMessage, validatorOptions); const errors = await validate(miningSubmitMessage, validatorOptions);
@@ -308,7 +318,7 @@ export class StratumV1Client {
} else { } else {
console.error('Mining Submit validation error'); console.log('Mining Submit validation error');
const err = new StratumErrorMessage( const err = new StratumErrorMessage(
miningSubmitMessage.id, miningSubmitMessage.id,
eStratumErrorCode.OtherUnknown, eStratumErrorCode.OtherUnknown,
@@ -335,7 +345,7 @@ export class StratumV1Client {
&& this.clientAuthorization != null && this.clientAuthorization != null
&& this.stratumInitialized == false) { && this.stratumInitialized == false) {
await this.initStratum(); this.initStratum();
} }
} }
@@ -343,6 +353,12 @@ export class StratumV1Client {
private async initStratum() { private async initStratum() {
this.stratumInitialized = true; this.stratumInitialized = true;
if (this.validateHeaderCompliance(this.clientSubscription.userAgent)) {
console.log(`Non compliant connection from userAgent: ${this.clientSubscription.userAgent}`);
await this.socket.end();
return;
}
switch (this.clientSubscription.userAgent) { switch (this.clientSubscription.userAgent) {
case 'cpuminer': { case 'cpuminer': {
this.sessionDifficulty = 0.1; this.sessionDifficulty = 0.1;
@@ -360,6 +376,9 @@ export class StratumV1Client {
this.stratumSubscription = this.stratumV1JobsService.newMiningJob$.subscribe(async (jobTemplate) => { this.stratumSubscription = this.stratumV1JobsService.newMiningJob$.subscribe(async (jobTemplate) => {
try { try {
if(jobTemplate.blockData.clearJobs){
this.miningSubmissionHashes.clear();
}
await this.sendNewMiningJob(jobTemplate); await this.sendNewMiningJob(jobTemplate);
} catch (e) { } catch (e) {
await this.socket.end(); await this.socket.end();
@@ -373,34 +392,40 @@ export class StratumV1Client {
}, 60 * 1000) }, 60 * 1000)
); );
this.backgroundWork.push( // this.backgroundWork.push(
setInterval(async () => { // setInterval(async () => {
await this.statistics.saveShares(this.entity); // await this.statistics.saveShares(this.clientEntity);
}, 60 * 1000) // }, 60 * 1000)
); // );
} }
private async sendNewMiningJob(jobTemplate: IJobTemplate) { private async sendNewMiningJob(jobTemplate: IJobTemplate) {
let payoutInformation; let payoutInformation= [
const devFeeAddress = this.configService.get('DEV_FEE_ADDRESS'); { address: this.clientAuthorization.address, percent: 100 }
//50Th/s ];
this.noFee = false; // const devFeeAddress = this.configService.get('DEV_FEE_ADDRESS');
if (this.entity) { // //50Th/s
this.hashRate = await this.clientStatisticsService.getHashRateForSession(this.clientAuthorization.address, this.clientAuthorization.worker, this.extraNonceAndSessionId); // this.noFee = false;
this.noFee = this.hashRate != 0 && this.hashRate < 50000000000000; // if (this.clientEntity) {
} // this.hashRate = await this.clientStatisticsService.getHashRateForSession(this.clientEntity.id);
if (this.noFee || devFeeAddress == null || devFeeAddress.length < 1) { // // 250Gh/s
payoutInformation = [ // if(this.hashRate < 250000000000){
{ address: this.clientAuthorization.address, percent: 100 } // this.statistics.targetSubmitShareEveryNSeconds = 10;
]; // }
// this.noFee = this.hashRate != 0 && this.hashRate < 50000000000000;
// }
// if (this.noFee || devFeeAddress == null || devFeeAddress.length < 1) {
// payoutInformation = [
// { address: this.clientAuthorization.address, percent: 100 }
// ];
} else { // } else {
payoutInformation = [ // payoutInformation = [
{ address: devFeeAddress, percent: 1.5 }, // { address: devFeeAddress, percent: 1.5 },
{ address: this.clientAuthorization.address, percent: 98.5 } // { address: this.clientAuthorization.address, percent: 98.5 }
]; // ];
} // }
const networkConfig = this.configService.get('NETWORK'); const networkConfig = this.configService.get('NETWORK');
let network; let network;
@@ -438,11 +463,11 @@ export class StratumV1Client {
private async handleMiningSubmission(submission: MiningSubmitMessage) { private async handleMiningSubmission(submission: MiningSubmitMessage) {
if (this.entity == null) { if (this.clientEntity == null) {
if (this.creatingEntity == null) { if (this.creatingEntity == null) {
this.creatingEntity = new Promise(async (resolve, reject) => { this.creatingEntity = new Promise(async (resolve, reject) => {
try { try {
this.entity = await this.clientService.insert({ this.clientEntity = await this.clientService.insert({
sessionId: this.extraNonceAndSessionId, sessionId: this.extraNonceAndSessionId,
address: this.clientAuthorization.address, address: this.clientAuthorization.address,
clientName: this.clientAuthorization.worker, clientName: this.clientAuthorization.worker,
@@ -462,10 +487,24 @@ export class StratumV1Client {
} }
} }
const submissionHash = submission.hash();
if(this.miningSubmissionHashes.has(submissionHash)){
const err = new StratumErrorMessage(
submission.id,
eStratumErrorCode.DuplicateShare,
'Duplicate share').response();
const success = await this.write(err);
if (!success) {
return false;
}
return false;
}else{
this.miningSubmissionHashes.add(submissionHash);
}
const job = this.stratumV1JobsService.getJobById(submission.jobId); 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 // a miner may submit a job that doesn't exist anymore if it was removed by a new block notification (or expired, 5 min)
if (job == null) { if (job == null) {
const err = new StratumErrorMessage( const err = new StratumErrorMessage(
submission.id, submission.id,
@@ -478,8 +517,23 @@ export class StratumV1Client {
} }
return false; return false;
} }
const jobTemplate = this.stratumV1JobsService.getJobTemplateById(job.jobTemplateId); const jobTemplate = this.stratumV1JobsService.getJobTemplateById(job.jobTemplateId);
if (jobTemplate == null) {
const err = new StratumErrorMessage(
submission.id,
eStratumErrorCode.JobNotFound,
'Job Template not found').response();
//console.log(err);
const success = await this.write(err);
if (!success) {
return false;
}
return false;
}
const updatedJobBlock = job.copyAndUpdateBlock( const updatedJobBlock = job.copyAndUpdateBlock(
jobTemplate, jobTemplate,
parseInt(submission.versionMask, 16), parseInt(submission.versionMask, 16),
@@ -515,33 +569,23 @@ export class StratumV1Client {
} }
} }
try { try {
await this.statistics.addShares(this.sessionDifficulty); await this.statistics.addShares(this.clientEntity, this.sessionDifficulty);
const now = new Date(); const now = new Date();
// only update every minute // only update every minute
if (this.entity.updatedAt == null || now.getTime() - this.entity.updatedAt.getTime() > 1000 * 60) { //if (this.clientEntity.updatedAt == null || now.getTime() - this.clientEntity.updatedAt.getTime() > 1000 * 60) {
await this.clientService.heartbeat(this.entity.address, this.entity.clientName, this.entity.sessionId, this.hashRate, now); await this.clientService.heartbeatBulkAsync(this.clientEntity.id, this.statistics.hashRate, now);
this.entity.updatedAt = now; this.clientEntity.updatedAt = now;
} //}
} catch (e) { } catch (e) {
console.log(e); console.log(e);
const err = new StratumErrorMessage(
submission.id,
eStratumErrorCode.DuplicateShare,
'Duplicate share').response();
console.error(err);
const success = await this.write(err);
if (!success) {
return false;
}
return false;
} }
if (submissionDifficulty > this.entity.bestDifficulty) { if (submissionDifficulty > this.clientEntity.bestDifficulty) {
await this.clientService.updateBestDifficulty(this.extraNonceAndSessionId, submissionDifficulty); await this.clientService.updateBestDifficulty(this.clientEntity.id, submissionDifficulty);
this.entity.bestDifficulty = submissionDifficulty; this.clientEntity.bestDifficulty = submissionDifficulty;
if (submissionDifficulty > (await this.addressSettingsService.getSettings(this.clientAuthorization.address, true)).bestDifficulty) { if (submissionDifficulty > (await this.addressSettingsService.getSettings(this.clientAuthorization.address, true)).bestDifficulty) {
await this.addressSettingsService.updateBestDifficulty(this.clientAuthorization.address, submissionDifficulty); await this.addressSettingsService.updateBestDifficulty(this.clientAuthorization.address, submissionDifficulty, this.clientEntity.userAgent);
} }
} }
@@ -585,9 +629,9 @@ export class StratumV1Client {
await this.socket.write(data); await this.socket.write(data);
// we need to clear the jobs so that the difficulty set takes effect. Otherwise the different miner implementations can cause issues
const jobTemplate = await firstValueFrom(this.stratumV1JobsService.newMiningJob$); const jobTemplate = await firstValueFrom(this.stratumV1JobsService.newMiningJob$);
// we need to clear the jobs so that the difficulty set takes effect. Otherwise the different miner implementations can cause issues
jobTemplate.blockData.clearJobs = true;
await this.sendNewMiningJob(jobTemplate); await this.sendNewMiningJob(jobTemplate);
} }
@@ -615,6 +659,18 @@ export class StratumV1Client {
return number; return number;
} }
private validateHeaderCompliance(userAgent: string): boolean {
const headerCompliance = this.configService.get<string>('COMPLIANT_HEADERS');
if (!headerCompliance || headerCompliance.trim() === '') {
return false;
}
const complianceList = headerCompliance.split(',').map(ua => ua.trim().toLowerCase());
const userAgentLower = userAgent.toLowerCase();
return complianceList.some(compliant => compliant.length > 0 && userAgentLower.includes(compliant));
}
private async write(message: string): Promise<boolean> { private async write(message: string): Promise<boolean> {
try { try {
if (!this.socket.destroyed && !this.socket.writableEnded) { if (!this.socket.destroyed && !this.socket.writableEnded) {
@@ -631,7 +687,7 @@ export class StratumV1Client {
return true; return true;
} else { } else {
console.error(`Error: Cannot write to closed or ended socket. ${this.extraNonceAndSessionId} ${message}`); //console.error(`Error: Cannot write to closed or ended socket. ${this.extraNonceAndSessionId} ${message}`);
this.destroy(); this.destroy();
if (!this.socket.destroyed) { if (!this.socket.destroyed) {
this.socket.destroy(); this.socket.destroy();
@@ -645,7 +701,7 @@ export class StratumV1Client {
} else if (!this.socket.destroyed) { } else if (!this.socket.destroyed) {
this.socket.destroy(); this.socket.destroy();
} }
console.error(`Error occurred while writing to socket: ${this.extraNonceAndSessionId}`, error); //console.error(`Error occurred while writing to socket: ${this.extraNonceAndSessionId}`, error);
return false; return false;
} }
} }
+73 -30
View File
@@ -2,66 +2,109 @@ import { ClientStatisticsService } from '../ORM/client-statistics/client-statist
import { ClientEntity } from '../ORM/client/client.entity'; import { ClientEntity } from '../ORM/client/client.entity';
const CACHE_SIZE = 30; const CACHE_SIZE = 30;
const TARGET_SUBMISSION_PER_SECOND = 10; const MIN_DIFF = 0.001;
const MIN_DIFF = 0.00001;
export class StratumV1ClientStatistics { export class StratumV1ClientStatistics {
private shareBacklog: number = 0; public targetSubmitShareEveryNSeconds: number = 30;
public hashRate = 0;
private shares: number = 0;
private acceptedCount: number = 0;
private submissionCacheStart: Date; private submissionCacheStart: Date;
private submissionCache = []; private submissionCache: { time: Date, difficulty: number }[] = [];
private submissionCacheDifficultySum = 0;
private currentTimeSlot: number = null;
constructor( constructor(
private readonly clientStatisticsService: ClientStatisticsService private readonly clientStatisticsService: ClientStatisticsService,
) { ) {
this.submissionCacheStart = new Date(); this.submissionCacheStart = new Date();
} }
public async saveShares(client: ClientEntity) {
if (client == null || client.address == null || client.clientName == null || client.sessionId == null) { // We don't want to save them here because it can be DB intensive, instead do it every once in
return; // awhile with saveShares()
} public async addShares(client: ClientEntity, targetDifficulty: number) {
// 10 min // 10 min
var coeff = 1000 * 60 * 10; var coeff = 1000 * 60 * 10;
var date = new Date(); var date = new Date();
var rounded = new Date(Math.floor(date.getTime() / coeff) * coeff); var timeSlot = new Date(Math.floor(date.getTime() / coeff) * coeff).getTime();
await this.clientStatisticsService.save({
time: rounded.getTime(),
shares: this.shareBacklog,
address: client.address,
clientName: client.clientName,
sessionId: client.sessionId
});
this.shareBacklog = 0;
}
// We don't want to save them here because it can be DB intensive, stead do it every once in
// awhile with saveShares()
public async addShares(targetDifficulty: number) {
if (this.submissionCache.length > CACHE_SIZE) { if (this.submissionCache.length > CACHE_SIZE) {
this.submissionCacheDifficultySum -= this.submissionCache[0].difficulty;
this.submissionCache.shift(); this.submissionCache.shift();
} }
this.submissionCache.push({ this.submissionCache.push({
time: new Date(), time: date,
difficulty: targetDifficulty, difficulty: targetDifficulty,
}); });
this.submissionCacheDifficultySum += targetDifficulty;
this.shareBacklog += targetDifficulty; if (this.currentTimeSlot == null) {
// First record, insert it
this.currentTimeSlot = timeSlot;
this.shares += targetDifficulty;
this.acceptedCount++;
await this.clientStatisticsService.insert({
time: this.currentTimeSlot,
clientId: client.id,
shares: this.shares,
acceptedCount: this.acceptedCount,
address: client.address,
clientName: client.clientName,
sessionId: client.sessionId
});
} else if (this.currentTimeSlot != timeSlot) {
// Transitioning to a new time slot,
// First update the old time slot with the latest data
this.clientStatisticsService.updateBulkAsync({
time: this.currentTimeSlot,
clientId: client.id,
shares: this.shares,
acceptedCount: this.acceptedCount,
});
// Set the new time slot and add incoming shares then insert it
this.currentTimeSlot = timeSlot;
this.shares = targetDifficulty;
this.acceptedCount = 1
await this.clientStatisticsService.insert({
time: this.currentTimeSlot,
clientId: client.id,
shares: this.shares,
acceptedCount: this.acceptedCount,
address: client.address,
clientName: client.clientName,
sessionId: client.sessionId
});
} else {
// Accept the shares if none of the prior conditions are met,
// saving to memory for storing later
this.shares += targetDifficulty;
this.acceptedCount++;
this.clientStatisticsService.updateBulkAsync({
time: this.currentTimeSlot,
clientId: client.id,
shares: this.shares,
acceptedCount: this.acceptedCount,
});
}
const time = new Date().getTime() - this.submissionCache[0].time.getTime();
if(time > 60000 && this.submissionCache.length > 2) {
this.hashRate = (this.submissionCacheDifficultySum * 4294967296) / (time / 1000);
}
} }
public getSuggestedDifficulty(clientDifficulty: number) { public getSuggestedDifficulty(clientDifficulty: number) {
// miner hasn't submitted shares in one minute // miner hasn't submitted shares in one minute
if (this.submissionCache.length < 5) { if (this.submissionCache.length < 5) {
if ((new Date().getTime() - this.submissionCacheStart.getTime()) / 1000 > 60) { if ((new Date().getTime() - this.submissionCacheStart.getTime()) / 5000 > 60) {
return this.nearestPowerOfTwo(clientDifficulty / 6); return this.nearestPowerOfTwo(clientDifficulty / 6);
} else { } else {
return null; return null;
@@ -76,7 +119,7 @@ export class StratumV1ClientStatistics {
const difficultyPerSecond = sum / diffSeconds; const difficultyPerSecond = sum / diffSeconds;
const targetDifficulty = difficultyPerSecond * TARGET_SUBMISSION_PER_SECOND; const targetDifficulty = difficultyPerSecond * this.targetSubmitShareEveryNSeconds;
if ((clientDifficulty * 2) < targetDifficulty || (clientDifficulty / 2) > targetDifficulty) { if ((clientDifficulty * 2) < targetDifficulty || (clientDifficulty / 2) > targetDifficulty) {
return this.nearestPowerOfTwo(targetDifficulty) return this.nearestPowerOfTwo(targetDifficulty)
@@ -1,5 +1,5 @@
import { Expose, Transform } from 'class-transformer'; import { Expose, Transform } from 'class-transformer';
import { ArrayMaxSize, ArrayMinSize, IsArray, IsString, MaxLength } from 'class-validator'; import { ArrayMaxSize, ArrayMinSize, IsArray, IsNumber, IsOptional, IsString, MaxLength } from 'class-validator';
import { eRequestMethod } from '../enums/eRequestMethod'; import { eRequestMethod } from '../enums/eRequestMethod';
import { IsBitcoinAddress } from '../validators/bitcoin-address.validator'; import { IsBitcoinAddress } from '../validators/bitcoin-address.validator';
@@ -8,7 +8,7 @@ import { StratumBaseMessage } from './StratumBaseMessage';
export class AuthorizationMessage extends StratumBaseMessage { export class AuthorizationMessage extends StratumBaseMessage {
@IsArray() @IsArray()
@ArrayMinSize(2) @ArrayMinSize(1)
@ArrayMaxSize(2) @ArrayMaxSize(2)
params: string[]; params: string[];
@@ -24,18 +24,36 @@ export class AuthorizationMessage extends StratumBaseMessage {
@IsString() @IsString()
@MaxLength(64) @MaxLength(64)
@Transform(({ value, key, obj, type }) => { @Transform(({ value, key, obj, type }) => {
return obj.params[0].split('.')[1] == null ? 'worker' : obj.params[0].split('.')[1]; const workerName = obj.params[0].split('.')[1];
if (workerName == null) {
return 'worker';
}
return workerName;
}) })
public worker: string; public worker: string;
@Expose()
@IsNumber()
@Transform(({ value, key, obj, type }) => {
const password: string | null = obj.params[1];
if (password?.includes('d=')) {
return parseInt(password.split('d=')[1]);
}
return null;
})
@IsOptional()
public startingDiff: number;
@Expose() @Expose()
@IsString() @IsString()
@Transform(({ value, key, obj, type }) => { @Transform(({ value, key, obj, type }) => {
return obj.params[1]; return obj.params[1];
}) })
@MaxLength(64) @MaxLength(64)
public password: string; @IsOptional()
public password?: string;
constructor() { constructor() {
super(); super();
@@ -3,6 +3,8 @@ import { ArrayMaxSize, ArrayMinSize, IsArray, IsString } from 'class-validator';
import { eRequestMethod } from '../enums/eRequestMethod'; import { eRequestMethod } from '../enums/eRequestMethod';
import { StratumBaseMessage } from './StratumBaseMessage'; import { StratumBaseMessage } from './StratumBaseMessage';
import * as bitcoinjs from 'bitcoinjs-lib';
export class MiningSubmitMessage extends StratumBaseMessage { export class MiningSubmitMessage extends StratumBaseMessage {
@@ -63,6 +65,11 @@ export class MiningSubmitMessage extends StratumBaseMessage {
}; };
} }
public hash(): string{
const buffer = Buffer.from(this.versionMask + this.nonce + this.extraNonce2 + this.ntime + this.jobId);
return bitcoinjs.crypto.hash256(buffer).toString('base64');
}
@@ -1,10 +1,10 @@
import { IsEnum, IsNumber } from 'class-validator'; import { IsDefined, IsEnum, IsNumber, IsNumberString, isNumberString, IsOptional } from 'class-validator';
import { eRequestMethod } from '../enums/eRequestMethod'; import { eRequestMethod } from '../enums/eRequestMethod';
export class StratumBaseMessage { export class StratumBaseMessage {
@IsNumber() @IsDefined()
id?: number = null; id?: number | string = null;
@IsEnum(eRequestMethod) @IsEnum(eRequestMethod)
method: eRequestMethod; method: eRequestMethod;
} }
@@ -5,7 +5,7 @@ import { eStratumErrorCode } from '../enums/eStratumErrorCode';
export class StratumErrorMessage { export class StratumErrorMessage {
constructor( constructor(
private id: number = null, private id: number | string = null,
private errorCode: eStratumErrorCode, private errorCode: eStratumErrorCode,
private errorMessage: string, private errorMessage: string,
private validationErrors: ValidationError[] = [] private validationErrors: ValidationError[] = []
@@ -14,7 +14,7 @@ export class SubscriptionMessage extends StratumBaseMessage {
@IsString() @IsString()
@MaxLength(128) @MaxLength(128)
@Transform(({ value, key, obj, type }) => { @Transform(({ value, key, obj, type }) => {
return obj.params[0] == null ? 'unknown' : SubscriptionMessage.refineUserAgent(obj.params[0]); return obj?.params?.[0] == null ? 'unknown' : SubscriptionMessage.refineUserAgent(obj.params[0]);
}) })
public userAgent: string; public userAgent: string;
@@ -41,7 +41,7 @@ export class SubscriptionMessage extends StratumBaseMessage {
public static refineUserAgent(userAgent: string): string { public static refineUserAgent(userAgent: string): string {
// return userAgent; // return userAgent;
userAgent = userAgent.split(' ')[0].split('/')[0].split('V')[0]; userAgent = userAgent.split(' ')[0].split('/')[0].split('V')[0].split('-')[0];
if (userAgent.includes('bosminer') || userAgent.includes('bOS')) { if (userAgent.includes('bosminer') || userAgent.includes('bOS')) {
userAgent = 'Braiins OS'; userAgent = 'Braiins OS';
@@ -31,7 +31,7 @@ export class SuggestDifficulty extends StratumBaseMessage {
public response(difficulty: number) { public response(difficulty: number) {
return { return {
id: this.id, id: null,
method: eResponseMethod.SET_DIFFICULTY, method: eResponseMethod.SET_DIFFICULTY,
params: [difficulty] params: [difficulty]
} }
+86 -26
View File
@@ -1,9 +1,11 @@
import { Injectable, OnModuleInit } from '@nestjs/common'; import { Injectable, OnModuleInit } from '@nestjs/common';
import { Interval } from '@nestjs/schedule';
import { DataSource } from 'typeorm'; import { DataSource } from 'typeorm';
import { UserAgentReportService } from '../ORM/_views/user-agent-report/user-agent-report.service';
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service'; import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
import { ClientService } from '../ORM/client/client.service'; import { ClientService } from '../ORM/client/client.service';
import { HomeGraphService } from '../ORM/home-graph/home-graph.service';
import { RpcBlockService } from '../ORM/rpc-block/rpc-block.service';
@Injectable() @Injectable()
export class AppService implements OnModuleInit { export class AppService implements OnModuleInit {
@@ -11,42 +13,100 @@ export class AppService implements OnModuleInit {
constructor( constructor(
private readonly clientStatisticsService: ClientStatisticsService, private readonly clientStatisticsService: ClientStatisticsService,
private readonly clientService: ClientService, private readonly clientService: ClientService,
private readonly dataSource: DataSource private readonly rpcBlockService: RpcBlockService,
private readonly homeGraphService: HomeGraphService,
private readonly dataSource: DataSource,
private readonly userAgentReportService: UserAgentReportService
) { ) {
} }
async onModuleInit() { async onModuleInit() {
// if (process.env.NODE_APP_INSTANCE == '0') { if (process.env.MASTER == 'true') {
// await this.dataSource.query(`VACUUM;`);
// } setInterval(async () => {
await this.deleteOldStatistics();
}, 1000 * 60 * 60);
setInterval(async () => {
console.log('Killing dead clients');
while (await this.clientService.killDeadClients()) {
}
console.log('Finished killing clients');
}, 1000 * 60 * 5);
setInterval(async () => {
console.log('Deleting Old Blocks');
await this.rpcBlockService.deleteOldBlocks();
}, 1000 * 60 * 60 * 24);
setInterval(async () => {
await this.updateChart();
}, 1000 * 60 * 10);
setInterval(async () => {
console.log('Refreshing user agent report view')
await this.userAgentReportService.refreshReport();
console.log('Finished Refreshing user agent report view')
}, 1000 * 60 * 5);
}
setInterval(async () => {
//console.log('Bulk update client stats');
await this.clientStatisticsService.doBulkAsyncUpdate();
}, 1000 * 30);
setInterval(async () => {
//console.log('Bulk update client stats');
await this.clientService.doBulkHeartbeatUpdate();
}, 1000 * 30);
//https://phiresky.github.io/blog/2020/sqlite-performance-tuning/
// //500 MB DB cache
// await this.dataSource.query(`PRAGMA cache_size = -500000;`);
//Normal is still completely corruption safe in WAL mode, and means only WAL checkpoints have to wait for FSYNC.
await this.dataSource.query(`PRAGMA synchronous = off;`);
// //6Gb
// await this.dataSource.query(`PRAGMA mmap_size = 6000000000;`);
} }
@Interval(1000 * 60 * 60)
private async deleteOldStatistics() { private async deleteOldStatistics() {
console.log('Deleting statistics'); 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(); const deletedStatistics = await this.clientStatisticsService.deleteOldStatistics();
console.log(`Deleted ${deletedStatistics.affected} old statistics`); console.log(`Deleted ${deletedStatistics.affected} old statistics`);
const deletedClients = await this.clientService.deleteOldClients(); const deletedClients = await this.clientService.deleteOldClients();
console.log(`Deleted ${deletedClients.affected} old clients`); 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();
}
}
private async updateChart() {
console.log('Updating Chart');
const latestGraphUpdate = await this.homeGraphService.getLatestTime();
const data = await this.dataSource.query(`
SELECT
time AS label,
ROUND(((SUM(shares) * 4294967296) / 600)) AS data
FROM
client_statistics_entity AS entry
WHERE
entry.time > ${latestGraphUpdate.getTime()}
GROUP BY
time
ORDER BY
time
LIMIT 144;
`)
console.log(`Fetched ${data.length} rows`);
if (data.length < 2) {
return;
}
const result = data.slice(0, data.length - 1);
this.homeGraphService.save(result);
}
} }
+88 -74
View File
@@ -1,26 +1,38 @@
import { Injectable } from '@nestjs/common'; import { Injectable, OnModuleInit } from '@nestjs/common';
import { ConfigService } from '@nestjs/config'; import { ConfigService } from '@nestjs/config';
import { RPCClient } from 'rpc-bitcoin'; 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 { 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 { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo'; import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo';
import * as PGPubsub from 'pg-pubsub';
@Injectable() @Injectable()
export class BitcoinRpcService { export class BitcoinRpcService implements OnModuleInit {
private blockHeight = 0;
private client: RPCClient; private client: RPCClient;
private _newBlock$: BehaviorSubject<IMiningInfo> = new BehaviorSubject(undefined); private _newBlockTemplate$: BehaviorSubject<IBlockTemplate> = new BehaviorSubject(undefined);
public newBlock$ = this._newBlock$.pipe(filter(block => block != null), shareReplay({ refCount: true, bufferSize: 1 })); private pubsubInstance: PGPubsub;
private resetTemplateInterval$ = new Subject<void>();
public miningInfo: IMiningInfo;
public newBlockTemplate$ = this._newBlockTemplate$.pipe(filter(block => block != null), shareReplay({ refCount: true, bufferSize: 1 }));
constructor( constructor(
private readonly configService: ConfigService, private readonly configService: ConfigService,
private rpcBlockService: RpcBlockService 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 url = this.configService.get('BITCOIN_RPC_URL');
const user = this.configService.get('BITCOIN_RPC_USER'); const user = this.configService.get('BITCOIN_RPC_USER');
const pass = this.configService.get('BITCOIN_RPC_PASSWORD'); const pass = this.configService.get('BITCOIN_RPC_PASSWORD');
@@ -34,85 +46,87 @@ export class BitcoinRpcService {
}, () => { }, () => {
console.error('Could not reach RPC host'); 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}`)
const sock = zmq.socket("sub"); if (process.env.MASTER != 'true') {
sock.connect(this.configService.get('BITCOIN_ZMQ_HOST')); this.pubsubInstance.addChannel('miningInfo', async (miningInfo: IMiningInfo) => {
sock.subscribe("rawblock"); //console.log('PG Sub. new template');
sock.on("message", async (topic: Buffer, message: Buffer) => { this.miningInfo = miningInfo;
console.log("new block zmq"); const savedBlockTemplate = await this.rpcBlockService.getSavedBlockTemplate(miningInfo.blocks);
await this.pollMiningInfo(); this._newBlockTemplate$.next(JSON.parse(savedBlockTemplate.data));
}); });
this.pollMiningInfo().then(() => { });
} else { } else {
setInterval(this.pollMiningInfo.bind(this), 500); console.log('Using ZMQ');
const sock = new zmq.Subscriber;
sock.connectTimeout = 1000;
sock.events.on('connect', () => {
console.log('ZMQ Connected');
});
sock.events.on('connect:retry', () => {
console.error('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);
// 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");
this.miningInfo = await this.getMiningInfo();
await this.getAndBroadcastLatestTemplate();
//Reset the block update interval
this.resetTemplateInterval$.next();
} }
} }
public async pollMiningInfo() { public async getAndBroadcastLatestTemplate() {
const miningInfo = await this.getMiningInfo(); const blockTemplate = await this.loadBlockTemplate(this.miningInfo.blocks);
if (miningInfo != null && miningInfo.blocks > this.blockHeight) { this._newBlockTemplate$.next(blockTemplate);
console.log("block height change"); await this.pubsubInstance.publish('miningInfo', this.miningInfo);
this._newBlock$.next(miningInfo);
this.blockHeight = miningInfo.blocks;
}
} }
private async waitForBlock(blockHeight: number): Promise<IBlockTemplate> { private async loadBlockTemplate(blockHeight: number) {
while (true) {
await new Promise(r => setTimeout(r, 100));
const block = await this.rpcBlockService.getBlock(blockHeight); console.log(`Master fetching block template ${blockHeight}`);
if (block != null && block.data != null) {
console.log('promise loop resolved');
return Promise.resolve(JSON.parse(block.data));
}
console.log('promise loop');
}
}
public async getBlockTemplate(blockHeight: number): Promise<IBlockTemplate> { let blockTemplate: IBlockTemplate;
let result: IBlockTemplate; while (blockTemplate == null) {
try { blockTemplate = await this.client.getblocktemplate({
template_request: {
const block = await this.rpcBlockService.getBlock(blockHeight); rules: ['segwit'],
mode: 'template',
if (block != null && block.data != null) { capabilities: ['serverlist', 'proposal']
return Promise.resolve(JSON.parse(block.data));
} else if (block == null) {
if (process.env.NODE_APP_INSTANCE != null) {
// There is a unique constraint on the block height so if another process tries to lock, it'll throw
try {
await this.rpcBlockService.lockBlock(blockHeight, process.env.NODE_APP_INSTANCE);
} catch (e) {
result = await this.waitForBlock(blockHeight);
}
} }
});
result = await this.client.getblocktemplate({
template_request: {
rules: ['segwit'],
mode: 'template',
capabilities: ['serverlist', 'proposal']
}
});
await this.rpcBlockService.saveBlock(blockHeight, JSON.stringify(result));
} else {
//wait for block
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; try {
console.log(`Saving block ${blockHeight}`);
await this.rpcBlockService.saveBlock(blockHeight, JSON.stringify(blockTemplate));
console.log('block saved');
} catch (e) {
console.error('Error saving block', e);
}
return blockTemplate;
} }
public async getMiningInfo(): Promise<IMiningInfo> { public async getMiningInfo(): Promise<IMiningInfo> {
+5 -5
View File
@@ -34,7 +34,7 @@ export class DiscordService implements OnModuleInit {
constructor(private readonly configService: ConfigService) { constructor(private readonly configService: ConfigService) {
if (process.env.NODE_APP_INSTANCE == null || process.env.NODE_APP_INSTANCE == '0') { if (process.env.MASTER == 'true') {
this.token = this.configService.get('DISCORD_BOT_TOKEN'); this.token = this.configService.get('DISCORD_BOT_TOKEN');
this.clientId = this.configService.get('DISCORD_BOT_CLIENTID'); this.clientId = this.configService.get('DISCORD_BOT_CLIENTID');
this.guildId = this.configService.get('DISCORD_BOT_GUILD_ID'); this.guildId = this.configService.get('DISCORD_BOT_GUILD_ID');
@@ -61,7 +61,7 @@ export class DiscordService implements OnModuleInit {
async onModuleInit(): Promise<void> { async onModuleInit(): Promise<void> {
if (process.env.NODE_APP_INSTANCE == null || process.env.NODE_APP_INSTANCE == '0') { if (process.env.MASTER == 'true') {
if (this.bot == null) { if (this.bot == null) {
return; return;
} }
@@ -93,7 +93,7 @@ export class DiscordService implements OnModuleInit {
} }
private async registerCommands() { private async registerCommands() {
if (process.env.NODE_APP_INSTANCE == null || process.env.NODE_APP_INSTANCE == '0') { if (process.env.MASTER == 'true') {
const rest = new REST().setToken(this.token); const rest = new REST().setToken(this.token);
try { try {
console.log(`Started refreshing ${commands.length} application (/) commands.`); console.log(`Started refreshing ${commands.length} application (/) commands.`);
@@ -113,7 +113,7 @@ export class DiscordService implements OnModuleInit {
} }
public async notifyRestarted() { public async notifyRestarted() {
if (process.env.NODE_APP_INSTANCE == null || process.env.NODE_APP_INSTANCE == '0') { if (process.env.MASTER == 'true') {
if (this.bot == null) { if (this.bot == null) {
return; return;
} }
@@ -125,7 +125,7 @@ export class DiscordService implements OnModuleInit {
} }
public async notifySubscribersBlockFound(height: number, block: Block, message: string) { public async notifySubscribersBlockFound(height: number, block: Block, message: string) {
if (process.env.NODE_APP_INSTANCE == null || process.env.NODE_APP_INSTANCE == '0') { if (process.env.MASTER == 'true') {
if (this.bot == null) { if (this.bot == null) {
return; return;
} }
-29
View File
@@ -1,29 +0,0 @@
import { Injectable, OnModuleInit } from '@nestjs/common';
import { Server, Socket } from 'net';
import { ProxyClient } from '../models/ProxyClient';
@Injectable()
export class ProxyService implements OnModuleInit {
constructor() {
}
async onModuleInit(): Promise<void> {
if (process.env.ENABLE_PROXY == 'true') {
console.log('Connecting to braiins');
const proxyServer = new Server(async (socket: Socket) => {
const proxyClient = new ProxyClient(socket);
});
proxyServer.listen(process.env.PROXY_PORT, () => {
console.log(`Proxy server is listening on port ${process.env.PROXY_PORT}`);
});
}
}
}
+39 -29
View File
@@ -13,6 +13,7 @@ export interface IJobTemplate {
merkle_branch: string[]; merkle_branch: string[];
blockData: { blockData: {
id: string, id: string,
creation: number,
coinbasevalue: number; coinbasevalue: number;
networkDifficulty: number; networkDifficulty: number;
height: number; height: number;
@@ -23,56 +24,42 @@ export interface IJobTemplate {
@Injectable() @Injectable()
export class StratumV1JobsService { export class StratumV1JobsService {
private lastIntervalCount: number;
private skipNext: boolean = false;
public newMiningJob$: Observable<IJobTemplate>; public newMiningJob$: Observable<IJobTemplate>;
public latestJobId: number = 1; public latestJobId: number = 1;
public latestJobTemplateId: number = 1; public latestJobTemplateId: number = 1;
public jobs: { [jobId: string]: MiningJob } = {}; public jobs: { [jobId: string]: MiningJob } = {};
public blocks: { [id: number]: IJobTemplate } = {}; public blocks: { [id: number]: IJobTemplate } = {};
// offset the interval so that all the cluster processes don't try and refresh at the same time. private lastBlockHeight = 0;
private delay = process.env.NODE_APP_INSTANCE == null ? 0 : parseInt(process.env.NODE_APP_INSTANCE) * 5000;
constructor( constructor(
private readonly bitcoinRpcService: BitcoinRpcService private readonly bitcoinRpcService: BitcoinRpcService
) { ) {
this.newMiningJob$ = combineLatest([this.bitcoinRpcService.newBlock$, interval(60000).pipe(delay(this.delay), startWith(-1))]).pipe( this.newMiningJob$ = this.bitcoinRpcService.newBlockTemplate$.pipe(
switchMap(([miningInfo, interval]) => { map((blockTemplate) => {
return from(this.bitcoinRpcService.getBlockTemplate(miningInfo.blocks)).pipe(map((blockTemplate) => {
return { if (process.env.MASTER == 'true') {
blockTemplate, console.log('Updating block template');
interval }
}
}))
}),
map(({ blockTemplate, interval }) => {
let clearJobs = false; 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; clearJobs = true;
this.skipNext = true; this.lastBlockHeight = currentBlockHeight;
console.log('new block')
} }
if (this.skipNext == true && clearJobs == false) { const currentTime = Math.floor(new Date().getTime() / 1000);
this.skipNext = false;
return null;
}
this.lastIntervalCount = interval;
return { return {
version: blockTemplate.version, version: blockTemplate.version,
bits: parseInt(blockTemplate.bits, 16), bits: parseInt(blockTemplate.bits, 16),
prevHash: this.convertToLittleEndian(blockTemplate.previousblockhash), prevHash: this.convertToLittleEndian(blockTemplate.previousblockhash),
transactions: blockTemplate.transactions.map(t => bitcoinjs.Transaction.fromHex(t.data)), transactions: blockTemplate.transactions.map(t => bitcoinjs.Transaction.fromHex(t.data)),
coinbasevalue: blockTemplate.coinbasevalue, coinbasevalue: blockTemplate.coinbasevalue,
timestamp: Math.floor(new Date().getTime() / 1000), timestamp: blockTemplate.mintime > currentTime ? blockTemplate.mintime : currentTime,
networkDifficulty: this.calculateNetworkDifficulty(parseInt(blockTemplate.bits, 16)), networkDifficulty: this.calculateNetworkDifficulty(parseInt(blockTemplate.bits, 16)),
clearJobs, clearJobs,
height: blockTemplate.height height: blockTemplate.height
@@ -113,6 +100,7 @@ export class StratumV1JobsService {
merkle_branch, merkle_branch,
blockData: { blockData: {
id, id,
creation: new Date().getTime(),
coinbasevalue, coinbasevalue,
networkDifficulty, networkDifficulty,
height, height,
@@ -124,11 +112,32 @@ export class StratumV1JobsService {
if (data.blockData.clearJobs) { if (data.blockData.clearJobs) {
this.blocks = {}; this.blocks = {};
this.jobs = {}; this.jobs = {};
}else{
let templatesDeleted = 0;
let jobsDeleted = 0;
const now = new Date().getTime();
// Delete old templates (5 minutes)
for(const templateId in this.blocks){
if(now - this.blocks[templateId].blockData.creation > (1000 * 60 * 5)){
delete this.blocks[templateId];
templatesDeleted++;
}
}
// Delete old jobs (5 minutes)
for (const jobId in this.jobs) {
if(now - this.jobs[jobId].creation > (1000 * 60 * 5)){
delete this.jobs[jobId];
jobsDeleted++;
}
}
//console.log(`Deleted ${templatesDeleted} templates and ${jobsDeleted} jobs.`)
} }
this.blocks[data.blockData.id] = data; this.blocks[data.blockData.id] = data;
}), }),
shareReplay({ refCount: true, bufferSize: 1 }) shareReplay({ refCount: true, bufferSize: 1 })
) )
this.newMiningJob$.subscribe();
} }
private calculateNetworkDifficulty(nBits: number) { private calculateNetworkDifficulty(nBits: number) {
@@ -137,7 +146,8 @@ export class StratumV1JobsService {
const target: number = mantissa * Math.pow(256, (exponent - 3)); // Calculate the target value const target: number = mantissa * Math.pow(256, (exponent - 3)); // Calculate the target value
const difficulty: number = (Math.pow(2, 208) * 65535) / target; // Calculate the difficulty const maxTarget = Math.pow(2, 208) * 65535; // Easiest target (max_target)
const difficulty: number = maxTarget / target; // Calculate the difficulty
return difficulty; return difficulty;
} }
+179 -59
View File
@@ -11,81 +11,201 @@ import { BitcoinRpcService } from './bitcoin-rpc.service';
import { NotificationService } from './notification.service'; import { NotificationService } from './notification.service';
import { StratumV1JobsService } from './stratum-v1-jobs.service'; import { StratumV1JobsService } from './stratum-v1-jobs.service';
import { readFileSync } from 'fs';
import { TlsOptions, TLSSocket, createServer } from 'tls';
import * as path from 'path';
@Injectable() @Injectable()
export class StratumV1Service implements OnModuleInit { export class StratumV1Service implements OnModuleInit {
constructor( private socketTimeout = 0;
private readonly bitcoinRpcService: BitcoinRpcService, private emptySocket = 0;
private readonly clientService: ClientService, private normalClosure = 0;
private readonly clientStatisticsService: ClientStatisticsService, private errorClosure = 0;
private readonly notificationService: NotificationService,
private readonly blocksService: BlocksService,
private readonly configService: ConfigService,
private readonly stratumV1JobsService: StratumV1JobsService,
private readonly addressSettingsService: AddressSettingsService
) {
} constructor(
private readonly bitcoinRpcService: BitcoinRpcService,
async onModuleInit(): Promise<void> { private readonly clientService: ClientService,
private readonly clientStatisticsService: ClientStatisticsService,
if (process.env.ENABLE_SOLO == 'true') { private readonly notificationService: NotificationService,
//await this.clientStatisticsService.deleteAll(); private readonly blocksService: BlocksService,
if (process.env.NODE_APP_INSTANCE == '0') { private readonly configService: ConfigService,
await this.clientService.deleteAll(); private readonly stratumV1JobsService: StratumV1JobsService,
} private readonly addressSettingsService: AddressSettingsService
setTimeout(() => { ) {
this.startSocketServer();
}, 1000 * 10)
} }
}
private startSocketServer() { async onModuleInit(): Promise<void> {
const server = new Server(async (socket: Socket) => {
//5 min if (process.env.MASTER == 'true') {
socket.setTimeout(1000 * 60 * 5); await this.clientService.deleteAll();
const client = new StratumV1Client(
socket,
this.stratumV1JobsService,
this.bitcoinRpcService,
this.clientService,
this.clientStatisticsService,
this.notificationService,
this.blocksService,
this.configService,
this.addressSettingsService
);
socket.on('close', async (hadError: boolean) => {
if (client.extraNonceAndSessionId != null) {
// Handle socket disconnection
await client.destroy();
console.log(`Client ${client.extraNonceAndSessionId} disconnected, hadError?:${hadError}`);
} }
});
socket.on('timeout', () => { // wait for all the other processes to init for an even connection distribution
console.log('socket timeout'); setTimeout(() => {
socket.end(); process.env.STRATUM_PORTS.split(',').forEach(port => {
socket.destroy(); this.startSocketServer(parseInt(port));
}); });
if (process.env.STRATUM_SECURE?.toLowerCase() === 'true') {
process.env.SECURE_STRATUM_PORTS.split(',').forEach(port => {
this.startSecureSocketServer(parseInt(port));
});
}
}, (10000));
socket.on('error', async (error: Error) => { }); setInterval(() => {
console.log(`Socket stats: ${this.emptySocket} empty, ${this.socketTimeout} timeouts, ${this.normalClosure} normal closure, ${this.errorClosure} error closure`);
this.emptySocket = 0;
this.socketTimeout = 0;
this.normalClosure = 0;
this.errorClosure = 0;
}, 1000 * 60);
// //console.log(`Client disconnected, socket error, ${client.extraNonceAndSessionId}`); }
private startSocketServer(port: number) {
const server = new Server(async (socket: Socket) => {
// Set 15-minute timeout
socket.setTimeout(1000 * 60 * 15);
const client = new StratumV1Client(
socket,
this.stratumV1JobsService,
this.bitcoinRpcService,
this.clientService,
this.clientStatisticsService,
this.notificationService,
this.blocksService,
this.configService,
this.addressSettingsService
);
// Unified cleanup function
const cleanup = async (reason: string) => {
if (client.extraNonceAndSessionId != null) {
await client.destroy();
if (reason == 'Error') {
this.errorClosure++;
} else {
this.normalClosure++;
}
}
if (!socket.destroyed) {
socket.end();
socket.destroy();
}
};
// Handle client disconnection
socket.on('close', async (hadError: boolean) => {
await cleanup(hadError ? "Error" : "Normal Closure");
});
// Handle socket timeouts
socket.on('timeout', async () => {
if (socket.bytesRead == 0 || socket.bytesWritten == 0) {
this.emptySocket++;
} else {
this.socketTimeout++;
}
await cleanup("Timeout");
});
// Handle errors properly
socket.on('error', async (error: Error) => {
await cleanup("Error");
});
//
}); });
server.listen(process.env.STRATUM_PORT, () => { // Ensure server itself handles errors
console.log(`Stratum server is listening on port ${process.env.STRATUM_PORT}`); server.on('error', (err) => {
}); console.error(`Server error: ${err.message}`);
});
server.listen(port, () => {
console.log(`Stratum server is listening on port ${port}`);
});
}
private startSecureSocketServer(port: number) {
const currentDirectory = process.cwd();
const keyPath = path.join(currentDirectory, 'secrets', 'key.pem');
const certPath = path.join(currentDirectory, 'secrets', 'cert.pem');
const tlsOptions: TlsOptions = {
key: readFileSync(keyPath),
cert: readFileSync(certPath),
};
const server = createServer(tlsOptions, async (socket: TLSSocket) => {
// Set 15-minute timeout
socket.setTimeout(1000 * 60 * 15);
const client = new StratumV1Client(
socket,
this.stratumV1JobsService,
this.bitcoinRpcService,
this.clientService,
this.clientStatisticsService,
this.notificationService,
this.blocksService,
this.configService,
this.addressSettingsService
);
const cleanup = async (reason: string) => {
if (client.extraNonceAndSessionId != null) {
await client.destroy();
if (reason === 'Error') {
this.errorClosure++;
} else {
this.normalClosure++;
}
}
if (!socket.destroyed) {
socket.end();
socket.destroy();
}
};
socket.on('close', async (hadError: boolean) => {
await cleanup(hadError ? 'Error' : 'Normal Closure');
});
socket.on('timeout', async () => {
if (socket.bytesRead === 0 || socket.bytesWritten === 0) {
this.emptySocket++;
} else {
this.socketTimeout++;
}
await cleanup('Timeout');
});
socket.on('error', async (error: Error) => {
await cleanup('Error');
});
// your protocol handling stays the same
});
server.on('error', (err) => {
console.error(`Server error: ${err.message}`);
});
server.listen(port, () => {
console.log(`Stratum TLS server is listening on port ${port}`);
});
}
}
} }