Migrate pool storage to TimescaleDB and Redis

This commit is contained in:
Ben
2026-06-07 20:33:13 -04:00
parent fb367ac5f0
commit ec69f5fb6c
64 changed files with 4751 additions and 1146 deletions
@@ -0,0 +1,36 @@
import { MigrationInterface, QueryRunner } from 'typeorm';
export class ActiveOnlyUserAgentReport1780860200000 implements MigrationInterface {
name = 'ActiveOnlyUserAgentReport1780860200000';
public async up(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`DROP MATERIALIZED VIEW IF EXISTS "user_agent_report_view"`);
await queryRunner.query(`
CREATE MATERIALIZED VIEW "user_agent_report_view" AS
SELECT
"userAgent",
COUNT("userAgent") AS "count",
MAX("bestDifficulty") AS "bestDifficulty",
SUM("hashRate") AS "totalHashRate"
FROM "client_entity"
WHERE "deletedAt" IS NULL
GROUP BY "userAgent"
ORDER BY "totalHashRate" DESC
`);
}
public async down(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`DROP MATERIALIZED VIEW IF EXISTS "user_agent_report_view"`);
await queryRunner.query(`
CREATE MATERIALIZED VIEW "user_agent_report_view" AS
SELECT
"userAgent",
COUNT("userAgent") AS "count",
MAX("bestDifficulty") AS "bestDifficulty",
SUM("hashRate") AS "totalHashRate"
FROM "client_entity"
GROUP BY "userAgent"
ORDER BY "totalHashRate" DESC
`);
}
}
@@ -0,0 +1,215 @@
import { MigrationInterface, QueryRunner } from 'typeorm';
export class InitialTimescaleSchema1780859300000 implements MigrationInterface {
public name = 'InitialTimescaleSchema1780859300000';
public transaction = false;
public async up(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`CREATE EXTENSION IF NOT EXISTS "uuid-ossp"`);
await queryRunner.query(`CREATE EXTENSION IF NOT EXISTS timescaledb`);
await queryRunner.query(`
CREATE TABLE IF NOT EXISTS "address_settings_entity" (
"address" varchar(62) PRIMARY KEY,
"shares" integer NOT NULL DEFAULT 0,
"bestDifficulty" numeric NOT NULL DEFAULT 0,
"bestDifficultyUserAgent" varchar,
"miscCoinbaseScriptData" varchar,
"deletedAt" timestamptz,
"createdAt" timestamptz NOT NULL DEFAULT now(),
"updatedAt" timestamptz NOT NULL DEFAULT now()
)
`);
await queryRunner.query(`
CREATE TABLE IF NOT EXISTS "client_entity" (
"id" uuid PRIMARY KEY DEFAULT uuid_generate_v4(),
"address" varchar(62) NOT NULL,
"clientName" varchar(64) NOT NULL,
"sessionId" varchar(8) NOT NULL,
"userAgent" varchar(128),
"startTime" timestamptz NOT NULL,
"bestDifficulty" numeric NOT NULL DEFAULT 0,
"hashRate" numeric NOT NULL DEFAULT 0,
"deletedAt" timestamptz,
"createdAt" timestamptz NOT NULL DEFAULT now(),
"updatedAt" timestamptz NOT NULL DEFAULT now()
)
`);
await queryRunner.query(`CREATE INDEX IF NOT EXISTS "IDX_client_entity_address" ON "client_entity" ("address")`);
await queryRunner.query(`CREATE INDEX IF NOT EXISTS "idx_client_cleanup" ON "client_entity" ("id") WHERE "deletedAt" IS NULL`);
await queryRunner.query(`CREATE UNIQUE INDEX IF NOT EXISTS "IDX_unique_nonce" ON "client_entity" ("sessionId") WHERE "deletedAt" IS NOT NULL`);
await queryRunner.query(`
CREATE TABLE IF NOT EXISTS "accepted_share_entity" (
"id" uuid NOT NULL DEFAULT uuid_generate_v4(),
"acceptedAt" timestamptz NOT NULL,
"protocol" varchar(8) NOT NULL,
"address" varchar(62) NOT NULL,
"clientName" varchar NOT NULL,
"sessionId" varchar(8) NOT NULL,
"clientId" uuid NOT NULL,
"jobId" varchar NOT NULL,
"jobTemplateId" varchar NOT NULL,
"blockHeight" bigint NOT NULL,
"creditedDifficulty" numeric NOT NULL,
"submissionDifficulty" numeric NOT NULL,
"networkDifficulty" numeric NOT NULL,
"nonce" varchar NOT NULL,
"ntime" varchar NOT NULL,
"version" varchar NOT NULL,
"extraNonce2" varchar NOT NULL,
"isBlockCandidate" boolean NOT NULL DEFAULT false,
"blockSubmissionResult" text,
"createdAt" timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY ("id", "acceptedAt")
)
`);
await queryRunner.query(`
SELECT create_hypertable(
'"accepted_share_entity"',
'acceptedAt',
chunk_time_interval => INTERVAL '1 day',
if_not_exists => TRUE
)
`);
await queryRunner.query(`CREATE INDEX IF NOT EXISTS "IDX_accepted_share_accounting_lookup" ON "accepted_share_entity" ("address", "clientName", "acceptedAt" DESC)`);
await queryRunner.query(`CREATE INDEX IF NOT EXISTS "IDX_accepted_share_client_lookup" ON "accepted_share_entity" ("clientId", "acceptedAt" DESC)`);
await queryRunner.query(`
CREATE UNIQUE INDEX IF NOT EXISTS "IDX_accepted_share_unique_submission"
ON "accepted_share_entity" ("acceptedAt", "protocol", "sessionId", "jobId", "nonce", "ntime", "version", "extraNonce2")
`);
await queryRunner.query(`
CREATE TABLE IF NOT EXISTS "blocks_entity" (
"id" bigserial PRIMARY KEY,
"height" integer NOT NULL,
"minerAddress" varchar(62) NOT NULL,
"worker" varchar NOT NULL,
"sessionId" varchar(8) NOT NULL,
"blockData" varchar NOT NULL,
"deletedAt" timestamptz,
"createdAt" timestamptz NOT NULL DEFAULT now(),
"updatedAt" timestamptz NOT NULL DEFAULT now()
)
`);
await queryRunner.query(`
CREATE TABLE IF NOT EXISTS "rpc_block_entity" (
"blockHeight" integer PRIMARY KEY,
"data" varchar
)
`);
await queryRunner.query(`
CREATE TABLE IF NOT EXISTS "telegram_subscriptions_entity" (
"id" bigserial PRIMARY KEY,
"address" varchar(62) NOT NULL,
"telegramChatId" integer NOT NULL,
"deletedAt" timestamptz,
"createdAt" timestamptz NOT NULL DEFAULT now(),
"updatedAt" timestamptz NOT NULL DEFAULT now()
)
`);
await queryRunner.query(`CREATE INDEX IF NOT EXISTS "IDX_telegram_subscriptions_address" ON "telegram_subscriptions_entity" ("address")`);
await queryRunner.query(`
CREATE MATERIALIZED VIEW IF NOT EXISTS "user_agent_report_view" AS
SELECT
"userAgent",
COUNT("userAgent") AS "count",
MAX("bestDifficulty") AS "bestDifficulty",
SUM("hashRate") AS "totalHashRate"
FROM "client_entity"
WHERE "deletedAt" IS NULL
GROUP BY "userAgent"
ORDER BY "totalHashRate" DESC
`);
await queryRunner.query(`
CREATE MATERIALIZED VIEW IF NOT EXISTS "accepted_share_10m"
WITH (timescaledb.continuous) AS
SELECT
time_bucket(INTERVAL '10 minutes', "acceptedAt") AS "bucket",
"address",
"clientName",
"sessionId",
"clientId",
SUM("creditedDifficulty") AS "shares",
COUNT(*) AS "acceptedCount",
ROUND(((SUM("creditedDifficulty") * 4294967296) / 600)) AS "hashRate"
FROM "accepted_share_entity"
GROUP BY "bucket", "address", "clientName", "sessionId", "clientId"
WITH NO DATA
`);
await queryRunner.query(`
CREATE MATERIALIZED VIEW IF NOT EXISTS "accepted_share_1h"
WITH (timescaledb.continuous) AS
SELECT
time_bucket(INTERVAL '1 hour', "acceptedAt") AS "bucket",
"address",
"clientName",
SUM("creditedDifficulty") AS "shares",
COUNT(*) AS "acceptedCount"
FROM "accepted_share_entity"
GROUP BY "bucket", "address", "clientName"
WITH NO DATA
`);
await queryRunner.query(`
CREATE MATERIALIZED VIEW IF NOT EXISTS "accepted_share_1d"
WITH (timescaledb.continuous) AS
SELECT
time_bucket(INTERVAL '1 day', "acceptedAt") AS "bucket",
"address",
SUM("creditedDifficulty") AS "shares",
COUNT(*) AS "acceptedCount"
FROM "accepted_share_entity"
GROUP BY "bucket", "address"
WITH NO DATA
`);
await queryRunner.query(`
SELECT add_continuous_aggregate_policy(
'accepted_share_10m',
start_offset => INTERVAL '2 days',
end_offset => INTERVAL '1 minute',
schedule_interval => INTERVAL '1 minute',
if_not_exists => TRUE
)
`);
await queryRunner.query(`
SELECT add_continuous_aggregate_policy(
'accepted_share_1h',
start_offset => INTERVAL '30 days',
end_offset => INTERVAL '10 minutes',
schedule_interval => INTERVAL '10 minutes',
if_not_exists => TRUE
)
`);
await queryRunner.query(`
SELECT add_continuous_aggregate_policy(
'accepted_share_1d',
start_offset => INTERVAL '365 days',
end_offset => INTERVAL '1 hour',
schedule_interval => INTERVAL '1 hour',
if_not_exists => TRUE
)
`);
await queryRunner.query(`SELECT add_retention_policy('accepted_share_entity', INTERVAL '90 days', if_not_exists => TRUE)`);
}
public async down(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`DROP MATERIALIZED VIEW IF EXISTS "accepted_share_1d"`);
await queryRunner.query(`DROP MATERIALIZED VIEW IF EXISTS "accepted_share_1h"`);
await queryRunner.query(`DROP MATERIALIZED VIEW IF EXISTS "accepted_share_10m"`);
await queryRunner.query(`DROP MATERIALIZED VIEW IF EXISTS "user_agent_report_view"`);
await queryRunner.query(`DROP TABLE IF EXISTS "accepted_share_entity" CASCADE`);
await queryRunner.query(`DROP TABLE IF EXISTS "telegram_subscriptions_entity"`);
await queryRunner.query(`DROP TABLE IF EXISTS "rpc_block_entity"`);
await queryRunner.query(`DROP TABLE IF EXISTS "blocks_entity"`);
await queryRunner.query(`DROP TABLE IF EXISTS "client_entity"`);
await queryRunner.query(`DROP TABLE IF EXISTS "address_settings_entity"`);
}
}
@@ -0,0 +1,33 @@
import { MigrationInterface, QueryRunner } from 'typeorm';
export class TimescaleOperationalHardening1780861200000 implements MigrationInterface {
public name = 'TimescaleOperationalHardening1780861200000';
public transaction = false;
public async up(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`
ALTER TABLE "accepted_share_entity" SET (
timescaledb.compress,
timescaledb.compress_orderby = '"acceptedAt" DESC',
timescaledb.compress_segmentby = '"address","clientName"'
)
`);
await queryRunner.query(`
SELECT add_compression_policy(
'accepted_share_entity',
INTERVAL '1 day',
if_not_exists => TRUE
)
`);
await queryRunner.query(`SELECT remove_retention_policy('accepted_share_entity', if_exists => TRUE)`);
await queryRunner.query(`SELECT add_retention_policy('accepted_share_entity', INTERVAL '30 days', if_not_exists => TRUE)`);
}
public async down(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`SELECT remove_compression_policy('accepted_share_entity', if_exists => TRUE)`);
await queryRunner.query(`SELECT remove_retention_policy('accepted_share_entity', if_exists => TRUE)`);
await queryRunner.query(`SELECT add_retention_policy('accepted_share_entity', INTERVAL '90 days', if_not_exists => TRUE)`);
}
}
+3 -3
View File
@@ -4,12 +4,12 @@ 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`
`CREATE UNIQUE INDEX IF NOT EXISTS "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"`);
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_unique_nonce"`);
}
}
}
@@ -3,13 +3,14 @@ import { TypeOrmModule } from '@nestjs/typeorm';
import { UserAgentReportService } from './user-agent-report.service';
import { UserAgentReportView } from './user-agent-report.view';
import { ClientEntity } from '../../client/client.entity';
@Global()
@Module({
imports: [TypeOrmModule.forFeature([UserAgentReportView])],
imports: [TypeOrmModule.forFeature([UserAgentReportView, ClientEntity])],
providers: [UserAgentReportService],
exports: [TypeOrmModule, UserAgentReportService],
})
export class UserAgentReportModule { }
export class UserAgentReportModule { }
@@ -2,7 +2,10 @@ import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { ClientEntity } from '../../client/client.entity';
import { ShareAccountingService } from '../../share-accounting/share-accounting.service';
import { UserAgentReportView } from './user-agent-report.view';
import { RedisMessagingService } from '../../../services/redis-messaging.service';
@Injectable()
export class UserAgentReportService {
@@ -10,26 +13,64 @@ export class UserAgentReportService {
constructor(
@InjectRepository(UserAgentReportView)
private userAgentReport: Repository<UserAgentReportView>,
@InjectRepository(ClientEntity)
private clientRepository: Repository<ClientEntity>,
private redisMessagingService: RedisMessagingService,
private shareAccountingService: ShareAccountingService,
) {
}
public async getReport() {
return await this.userAgentReport.find();
const presences = await this.redisMessagingService.getAllClientPresence();
const shareSummaries = await this.shareAccountingService.getSessionSummaries(
presences.map(presence => presence.clientId),
);
const rows = new Map<string, {
userAgent: string;
count: number;
bestDifficulty: number;
totalHashRate: number;
}>();
presences.forEach(presence => {
const userAgent = presence.userAgent == null || presence.userAgent.length === 0
? 'Other'
: presence.userAgent;
const row = rows.get(userAgent) ?? {
userAgent,
count: 0,
bestDifficulty: 0,
totalHashRate: 0,
};
const shareSummary = shareSummaries.get(presence.clientId);
row.count++;
row.bestDifficulty = Math.max(
row.bestDifficulty,
Number(presence.bestDifficulty ?? 0),
Number(shareSummary?.bestSubmissionDifficulty ?? 0),
);
row.totalHashRate += Number(shareSummary?.hashRateLast10Minutes ?? presence.hashRate ?? 0);
rows.set(userAgent, row);
});
return [...rows.values()]
.sort((left, right) => right.totalHashRate - left.totalHashRate)
.map(row => ({
userAgent: row.userAgent,
count: row.count.toString(),
bestDifficulty: row.bestDifficulty,
totalHashRate: row.totalHashRate.toString(),
}));
}
public async refreshReport() {
try {
return await this.userAgentReport.query(`
COMMIT;
BEGIN TRANSACTION ISOLATION LEVEL REPEATABLE READ;
REFRESH MATERIALIZED VIEW user_agent_report_view;
COMMIT;
`);
return await this.userAgentReport.query(`REFRESH MATERIALIZED VIEW user_agent_report_view`);
} catch (e) {
console.log(e)
}
}
}
}
@@ -12,6 +12,7 @@ import { ClientEntity } from '../../client/client.entity';
.addSelect('MAX(client.bestDifficulty)', 'bestDifficulty')
.addSelect('SUM(client.hashRate)', 'totalHashRate')
.from(ClientEntity, 'client')
.where('client.deletedAt IS NULL')
.groupBy('client.userAgent')
.orderBy('"totalHashRate"', 'DESC')
})
@@ -24,4 +25,4 @@ export class UserAgentReportView {
bestDifficulty: number;
@ViewColumn()
totalHashRate: string;
}
}
@@ -0,0 +1,67 @@
import { Column, Entity, Index, PrimaryColumn, PrimaryGeneratedColumn } from 'typeorm';
@Entity()
@Index('IDX_accepted_share_accounting_lookup', ['address', 'clientName', 'acceptedAt'])
@Index('IDX_accepted_share_client_lookup', ['clientId', 'acceptedAt'])
@Index('IDX_accepted_share_unique_submission', ['acceptedAt', 'protocol', 'sessionId', 'jobId', 'nonce', 'ntime', 'version', 'extraNonce2'], { unique: true })
export class AcceptedShareEntity {
@PrimaryGeneratedColumn('uuid')
id: string;
@PrimaryColumn({ type: 'timestamptz' })
acceptedAt: Date;
@Column({ length: 8, type: 'varchar' })
protocol: 'sv1' | 'sv2';
@Column({ length: 62, type: 'varchar' })
address: string;
@Column()
clientName: string;
@Column({ length: 8, type: 'varchar' })
sessionId: string;
@Column({ type: 'uuid' })
clientId: string;
@Column()
jobId: string;
@Column()
jobTemplateId: string;
@Column({ type: 'bigint' })
blockHeight: number;
@Column({ type: 'decimal' })
creditedDifficulty: number;
@Column({ type: 'decimal' })
submissionDifficulty: number;
@Column({ type: 'decimal' })
networkDifficulty: number;
@Column()
nonce: string;
@Column()
ntime: string;
@Column()
version: string;
@Column()
extraNonce2: string;
@Column({ default: false })
isBlockCandidate: boolean;
@Column({ type: 'text', nullable: true })
blockSubmissionResult?: string;
@Column({ type: 'timestamptz', default: () => 'CURRENT_TIMESTAMP' })
createdAt: Date;
}
@@ -1,48 +0,0 @@
import { Column, Entity, Index, ManyToOne } from 'typeorm';
import { ClientEntity } from '../client/client.entity';
import { PrimaryGeneratedBigIntColumn } from '../utils/PrimaryGeneratedBigIntColumn';
import { TrackedEntity } from '../utils/TrackedEntity.entity';
@Entity()
//Index for statistics save
@Index(["clientId", "time"])
export class ClientStatisticsEntity extends TrackedEntity {
@PrimaryGeneratedBigIntColumn()
id: number;
@Column({ length: 62, type: 'varchar' })
address: string;
@Column()
clientName: string;
@Column({ length: 8, type: 'varchar' })
sessionId: string;
@Index()
@Column({ type: 'bigint' })
time: number;
@Column({ type: 'decimal' })
shares: number;
@Column({ default: 0, type: 'bigint' })
acceptedCount: number;
@ManyToOne(
() => ClientEntity,
clientEntity => clientEntity.statistics,
{ nullable: false, }
)
client: ClientEntity;
@Index()
@Column({ name: 'clientId' })
public clientId: string;
}
@@ -1,14 +1,11 @@
import { Global, Module } from '@nestjs/common';
import { TypeOrmModule } from '@nestjs/typeorm';
import { ClientStatisticsEntity } from './client-statistics.entity';
import { ClientStatisticsService } from './client-statistics.service';
@Global()
@Module({
imports: [TypeOrmModule.forFeature([ClientStatisticsEntity])],
providers: [ClientStatisticsService],
exports: [TypeOrmModule, ClientStatisticsService],
exports: [ClientStatisticsService],
})
export class ClientStatisticsModule { }
export class ClientStatisticsModule { }
@@ -1,264 +1,129 @@
import { Injectable } from '@nestjs/common';
import { InjectDataSource, InjectRepository } from '@nestjs/typeorm';
import { DataSource, Repository } from 'typeorm';
import { ClientStatisticsEntity } from './client-statistics.entity';
import { InjectDataSource } from '@nestjs/typeorm';
import { DataSource } from 'typeorm';
const HASHES_PER_DIFFICULTY = 4294967296;
const CHART_BUCKET_SECONDS = 600;
const CHART_WINDOW = '24 hours';
const SITE_CHART_WINDOW = '7 days';
const REALTIME_WINDOW = '20 minutes';
@Injectable()
export class ClientStatisticsService {
private bulkAsyncUpdates: {
[key: string]: Partial<ClientStatisticsEntity>
} = {};
constructor(
@InjectDataSource()
private dataSource: DataSource,
@InjectRepository(ClientStatisticsEntity)
private clientStatisticsRepository: Repository<ClientStatisticsEntity>,
) {
}
// public async update(clientStatistic: Partial<ClientStatisticsEntity>) {
// await this.clientStatisticsRepository.update({ clientId: clientStatistic.clientId, time: clientStatistic.time },
// {
// shares: clientStatistic.shares,
// 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 getChartDataForSite(limit: number = 144 * 7) {
return this.getAcceptedShareChartData('', [], limit, SITE_CHART_WINDOW);
}
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() {
const oneDayAgo = new Date(Date.now() - 24 * 60 * 60 * 1000);
return await this.clientStatisticsRepository
.createQueryBuilder()
.delete()
.from(ClientStatisticsEntity)
.where('time < :time', { time: oneDayAgo.getTime() })
.execute();
}
public async getChartDataForAddress(address: string) {
var yesterday = new Date(new Date().getTime() - (24 * 60 * 60 * 1000));
const query = `
SELECT
time AS label,
(SUM(shares) * 4294967296) / 600 AS data
FROM
client_statistics_entity AS entry
WHERE
entry.address = $1 AND entry.time > $2
GROUP BY
time
ORDER BY
time
LIMIT 144;
`;
const result = await this.clientStatisticsRepository.query(query, [address, yesterday.getTime()]);
return result.map(res => {
res.label = new Date(parseInt(res.label)).toISOString();
return res;
}).slice(0, result.length - 1);
return this.getAcceptedShareChartData(
'AND "address" = $1',
[address],
144,
CHART_WINDOW,
);
}
public async getHashRateForGroup(address: string, clientName: string) {
var oneHour = new Date(new Date().getTime() - (60 * 60 * 1000));
const query = `
const result = await this.dataSource.query(`
SELECT
SUM(entry.shares) AS difficultySum
FROM
client_statistics_entity AS entry
WHERE
entry.address = $1 AND entry.clientName = $2 AND entry.time > ${oneHour.getTime()}
`;
const result = await this.clientStatisticsRepository.query(query, [address, clientName]);
const difficultySum = result[0].difficultySum;
return (difficultySum * 4294967296) / (600);
COALESCE((SUM("creditedDifficulty") * ${HASHES_PER_DIFFICULTY}) / ${CHART_BUCKET_SECONDS}, 0) AS "hashRate"
FROM "accepted_share_entity"
WHERE "address" = $1
AND "clientName" = $2
AND "acceptedAt" > NOW() - INTERVAL '1 hour'
`, [address, clientName]);
return parseFloat(result[0]?.hashRate ?? '0');
}
public async getChartDataForGroup(address: string, clientName: string) {
var yesterday = new Date(new Date().getTime() - (24 * 60 * 60 * 1000));
const query = `
SELECT
time AS label,
(SUM(shares) * 4294967296) / 600 AS data
FROM
client_statistics_entity AS entry
WHERE
entry.address = $1 AND entry."clientName" = $2 AND entry.time > ${yesterday.getTime()}
GROUP BY
time
ORDER BY
time
LIMIT 144;
`;
const result = await this.clientStatisticsRepository.query(query, [address, clientName]);
return result.map(res => {
res.label = new Date(parseInt(res.label)).toISOString();
return res;
}).slice(0, result.length - 1);
return this.getAcceptedShareChartData(
'AND "address" = $1 AND "clientName" = $2',
[address, clientName],
144,
CHART_WINDOW,
);
}
// public async getHashRateForSession(clientId: string) {
// const query = `
// SELECT
// "createdAt",
// "updatedAt",
// shares
// FROM
// client_statistics_entity AS entry
// WHERE
// entry."clientId" = $1
// ORDER BY time DESC
// LIMIT 2;
// `;
// const result = await this.clientStatisticsRepository.query(query, [clientId]);
// if (result.length < 1) {
// return 0;
// }
// const latestStat = result[0];
// if (result.length < 2) {
// const time = new Date(latestStat.updatedAt).getTime() - new Date(latestStat.createdAt).getTime();
// // 1min
// if (time < 1000 * 60) {
// return 0;
// }
// return (parseFloat(latestStat.shares) * 4294967296) / (time / 1000);
// } else {
// const secondLatestStat = result[1];
// const time = new Date(latestStat.updatedAt).getTime() - new Date(secondLatestStat.createdAt).getTime();
// // 1min
// if (time < 1000 * 60) {
// return 0;
// }
// return ((parseFloat(latestStat.shares) + parseFloat(secondLatestStat.shares)) * 4294967296) / (time / 1000);
// }
// }
public async getChartDataForSession(clientId: string) {
var yesterday = new Date(new Date().getTime() - (24 * 60 * 60 * 1000));
return this.getAcceptedShareChartData(
'AND "clientId" = $1',
[clientId],
144,
CHART_WINDOW,
);
}
private async getAcceptedShareChartData(filterSql: string, params: unknown[], limit: number, windowSql: string) {
const query = `
WITH bounds AS (
SELECT
NOW() - INTERVAL '${windowSql}' AS since,
time_bucket(INTERVAL '10 minutes', NOW() - INTERVAL '${REALTIME_WINDOW}') AS realtime_start
),
aggregate_rows AS (
SELECT
"bucket",
SUM("shares") AS "shares",
SUM("acceptedCount") AS "acceptedCount"
FROM "accepted_share_10m", bounds
WHERE "bucket" > bounds.since
AND "bucket" < bounds.realtime_start
${filterSql}
GROUP BY "bucket"
),
realtime_rows AS (
SELECT
time_bucket(INTERVAL '10 minutes', "acceptedAt") AS "bucket",
SUM("creditedDifficulty") AS "shares",
COUNT(*) AS "acceptedCount"
FROM "accepted_share_entity", bounds
WHERE "acceptedAt" > bounds.since
AND "acceptedAt" >= bounds.realtime_start
${filterSql}
GROUP BY "bucket"
),
combined_rows AS (
SELECT * FROM aggregate_rows
UNION ALL
SELECT * FROM realtime_rows
)
SELECT
time AS label,
(SUM(shares) * 4294967296) / 600 AS data
FROM
client_statistics_entity AS entry
WHERE
entry."clientId" = $1 AND entry.time > ${yesterday.getTime()}
GROUP BY
time
ORDER BY
time
LIMIT 144;
"label",
"data",
"shares",
"acceptedCount"
FROM (
SELECT
"bucket" AS "label",
ROUND((SUM("shares") * ${HASHES_PER_DIFFICULTY}) / ${CHART_BUCKET_SECONDS}) AS "data",
SUM("shares") AS "shares",
SUM("acceptedCount") AS "acceptedCount"
FROM combined_rows
GROUP BY "bucket"
ORDER BY "bucket" DESC
LIMIT ${limit}
) AS limited_rows
ORDER BY "label"
`;
const result = await this.clientStatisticsRepository.query(query, [clientId]);
const result = await this.dataSource.query(query, params);
return result.map(res => {
res.label = new Date(parseInt(res.label)).toISOString();
return res;
}).slice(0, result.length - 1);
}
public async deleteAll() {
return await this.clientStatisticsRepository.clear()
return {
label: new Date(res.label).toISOString(),
data: res.data,
shares: Number(res.shares ?? 0),
acceptedCount: Number(res.acceptedCount ?? 0),
};
});
}
}
+1 -8
View File
@@ -1,6 +1,5 @@
import { Column, Entity, Index, OneToMany, PrimaryGeneratedColumn } from 'typeorm';
import { Column, Entity, Index, PrimaryGeneratedColumn } from 'typeorm';
import { ClientStatisticsEntity } from '../client-statistics/client-statistics.entity';
import { TrackedEntity } from '../utils/TrackedEntity.entity';
@@ -39,10 +38,4 @@ export class ClientEntity extends TrackedEntity {
@Column({ default: 0, type: 'decimal' })
hashRate: number;
@OneToMany(
() => ClientStatisticsEntity,
clientStatisticsEntity => clientStatisticsEntity.client
)
statistics: ClientStatisticsEntity[]
}
+2 -95
View File
@@ -1,6 +1,6 @@
import { Injectable } from '@nestjs/common';
import { InjectDataSource, InjectRepository } from '@nestjs/typeorm';
import { DataSource, Repository } from 'typeorm';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { ClientEntity } from './client.entity';
@@ -9,106 +9,13 @@ import { ClientEntity } from './client.entity';
@Injectable()
export class ClientService {
private heartbeatBulkUpdate: { [id: string]: { id: string, hashRate: number, updatedAt: Date } } = {};
constructor(
@InjectDataSource()
private dataSource: DataSource,
@InjectRepository(ClientEntity)
private clientRepository: Repository<ClientEntity>
) {
}
// client.service.ts
public async killDeadClients(): Promise<boolean> {
const BATCH_SIZE = 10000;
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]);
if (deadClients.length === 0) {
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 heartbeat(id, hashRate: number, updatedAt: Date) {
// return await this.clientRepository.update({ id }, { hashRate, deletedAt: null, updatedAt });
// }
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> {
const insertResult = await this.clientRepository.insert(partialClient);
-16
View File
@@ -1,16 +0,0 @@
import { Column, Entity } from 'typeorm';
import { PrimaryGeneratedBigIntColumn } from '../utils/PrimaryGeneratedBigIntColumn';
@Entity()
export class HomeGraphEntity {
@PrimaryGeneratedBigIntColumn()
id: number;
@Column({ type: 'bigint' })
label: number;
@Column({ type: 'bigint' })
data: number;
}
-15
View File
@@ -1,15 +0,0 @@
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
@@ -1,45 +0,0 @@
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
}
});
}
}
@@ -0,0 +1,13 @@
import { Global, Module } from '@nestjs/common';
import { TypeOrmModule } from '@nestjs/typeorm';
import { AcceptedShareEntity } from '../accepted-share/accepted-share.entity';
import { ShareAccountingService } from './share-accounting.service';
@Global()
@Module({
imports: [TypeOrmModule.forFeature([AcceptedShareEntity])],
providers: [ShareAccountingService],
exports: [TypeOrmModule, ShareAccountingService],
})
export class ShareAccountingModule { }
@@ -0,0 +1,233 @@
import { ShareAccountingService } from './share-accounting.service';
describe('ShareAccountingService', () => {
const originalEnv = process.env;
beforeEach(() => {
jest.useRealTimers();
jest.restoreAllMocks();
process.env = { ...originalEnv };
});
afterAll(() => {
process.env = originalEnv;
});
it('should normalize and save accepted share records', async () => {
const repository = {
create: jest.fn(value => value),
insert: jest.fn().mockResolvedValue({}),
};
process.env.SHARE_ACCOUNTING_BATCH_SIZE = '1';
const service = new ShareAccountingService(repository as any);
const acceptedAt = new Date('2026-06-07T12:00:00Z');
await expect(service.recordAcceptedShare({
protocol: 'sv1',
acceptedAt,
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
clientName: 'worker',
sessionId: '57a6f098',
clientId: '6df7beea-921a-46cb-b7dc-e8075f916232',
jobId: '1',
jobTemplateId: '2',
blockHeight: 900000,
creditedDifficulty: 1024,
submissionDifficulty: 2048,
networkDifficulty: 100000,
nonce: 123,
ntime: 456,
version: 536870912,
extraNonce2: 'c708000000000000',
isBlockCandidate: false,
blockSubmissionResult: null,
})).resolves.toEqual(expect.objectContaining({
protocol: 'sv1',
creditedDifficulty: 1024,
}));
expect(repository.create).toHaveBeenCalledWith(expect.objectContaining({
acceptedAt,
nonce: '123',
ntime: '456',
version: '536870912',
blockSubmissionResult: null,
}));
expect(repository.insert).toHaveBeenCalledWith([expect.objectContaining({
protocol: 'sv1',
creditedDifficulty: 1024,
})]);
});
it('should batch accepted share inserts until the batch size is reached', async () => {
process.env.SHARE_ACCOUNTING_BATCH_SIZE = '2';
process.env.SHARE_ACCOUNTING_FLUSH_INTERVAL_MS = '1000';
const repository = {
create: jest.fn(value => value),
insert: jest.fn().mockResolvedValue({}),
};
const service = new ShareAccountingService(repository as any);
const first = service.recordAcceptedShare(buildRecord('1'));
expect(repository.insert).not.toHaveBeenCalled();
const second = service.recordAcceptedShare(buildRecord('2'));
await expect(Promise.all([first, second])).resolves.toHaveLength(2);
expect(repository.insert).toHaveBeenCalledTimes(1);
expect(repository.insert).toHaveBeenCalledWith([
expect.objectContaining({ jobId: '1' }),
expect.objectContaining({ jobId: '2' }),
]);
});
it('should flush partial batches on the flush interval', async () => {
jest.useFakeTimers();
process.env.SHARE_ACCOUNTING_BATCH_SIZE = '10';
process.env.SHARE_ACCOUNTING_FLUSH_INTERVAL_MS = '25';
const repository = {
create: jest.fn(value => value),
insert: jest.fn().mockResolvedValue({}),
};
const service = new ShareAccountingService(repository as any);
const pending = service.recordAcceptedShare(buildRecord('interval'));
expect(repository.insert).not.toHaveBeenCalled();
jest.advanceTimersByTime(25);
await expect(pending).resolves.toEqual(expect.objectContaining({ jobId: 'interval' }));
expect(repository.insert).toHaveBeenCalledTimes(1);
});
it('should reject accepted shares when the accounting queue is full', async () => {
process.env.SHARE_ACCOUNTING_BATCH_SIZE = '10';
process.env.SHARE_ACCOUNTING_FLUSH_INTERVAL_MS = '1000';
process.env.SHARE_ACCOUNTING_MAX_QUEUE_SIZE = '1';
const repository = {
create: jest.fn(value => value),
insert: jest.fn().mockResolvedValue({}),
};
const service = new ShareAccountingService(repository as any);
const queued = service.recordAcceptedShare(buildRecord('queued'));
await expect(service.recordAcceptedShare(buildRecord('overflow')))
.rejects
.toThrow('Share accounting queue is full');
await service.flushPendingShares();
await expect(queued).resolves.toEqual(expect.objectContaining({ jobId: 'queued' }));
});
it('should return numeric accounting summaries with protocol breakdown', async () => {
const repository = {
query: jest.fn()
.mockResolvedValueOnce([{
totalAcceptedShares: '3',
totalCreditedDifficulty: '96',
acceptedSharesLast10Minutes: '2',
creditedDifficultyLast10Minutes: '64',
acceptedSharesLastHour: '3',
creditedDifficultyLastHour: '96',
acceptedSharesLastDay: '3',
creditedDifficultyLastDay: '96',
hashRateLast10Minutes: '458129844.9',
hashRateLastHour: '114532461.2',
bestSubmissionDifficulty: '2048',
blockCandidateCount: '1',
latestShareAt: new Date('2026-06-07T12:10:00Z'),
}])
.mockResolvedValueOnce([{
protocol: 'sv1',
acceptedShares: '3',
creditedDifficulty: '96',
}]),
};
const service = new ShareAccountingService(repository as any);
await expect(service.getAddressSummary('bc1qtest')).resolves.toEqual({
totalAcceptedShares: 3,
totalCreditedDifficulty: 96,
acceptedSharesLast10Minutes: 2,
creditedDifficultyLast10Minutes: 64,
acceptedSharesLastHour: 3,
creditedDifficultyLastHour: 96,
acceptedSharesLastDay: 3,
creditedDifficultyLastDay: 96,
hashRateLast10Minutes: 458129844.9,
hashRateLastHour: 114532461.2,
bestSubmissionDifficulty: 2048,
blockCandidateCount: 1,
latestShareAt: '2026-06-07T12:10:00.000Z',
protocolBreakdown: [{
protocol: 'sv1',
acceptedShares: 3,
creditedDifficulty: 96,
}],
});
expect(repository.query).toHaveBeenNthCalledWith(
1,
expect.stringContaining('"address" = $1'),
['bc1qtest'],
);
expect(repository.query).toHaveBeenNthCalledWith(
2,
expect.stringContaining('GROUP BY "protocol"'),
['bc1qtest'],
);
});
it('should cache accounting summaries briefly to protect hot dashboard endpoints', async () => {
process.env.SHARE_ACCOUNTING_SUMMARY_CACHE_MS = '1000';
const repository = {
query: jest.fn()
.mockResolvedValueOnce([{
totalAcceptedShares: '1',
totalCreditedDifficulty: '32',
acceptedSharesLast10Minutes: '1',
creditedDifficultyLast10Minutes: '32',
acceptedSharesLastHour: '1',
creditedDifficultyLastHour: '32',
acceptedSharesLastDay: '1',
creditedDifficultyLastDay: '32',
hashRateLast10Minutes: '1',
hashRateLastHour: '1',
bestSubmissionDifficulty: '32',
blockCandidateCount: '0',
latestShareAt: null,
}])
.mockResolvedValueOnce([{ protocol: 'sv1', acceptedShares: '1', creditedDifficulty: '32' }]),
};
const service = new ShareAccountingService(repository as any);
await service.getPoolSummary();
await service.getPoolSummary();
expect(repository.query).toHaveBeenCalledTimes(2);
});
});
function buildRecord(jobId: string) {
return {
protocol: 'sv1' as const,
acceptedAt: new Date('2026-06-07T12:00:00Z'),
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
clientName: 'worker',
sessionId: '57a6f098',
clientId: '6df7beea-921a-46cb-b7dc-e8075f916232',
jobId,
jobTemplateId: 'template',
blockHeight: 900000,
creditedDifficulty: 1024,
submissionDifficulty: 2048,
networkDifficulty: 100000,
nonce: 123,
ntime: 456,
version: 536870912,
extraNonce2: 'c708000000000000',
isBlockCandidate: false,
blockSubmissionResult: null,
};
}
@@ -0,0 +1,374 @@
import { Injectable, OnModuleDestroy } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { AcceptedShareEntity } from '../accepted-share/accepted-share.entity';
export interface AcceptedShareRecord {
protocol: 'sv1' | 'sv2';
acceptedAt?: Date;
address: string;
clientName: string;
sessionId: string;
clientId: string;
jobId: string;
jobTemplateId: string;
blockHeight: number;
creditedDifficulty: number;
submissionDifficulty: number;
networkDifficulty: number;
nonce: string | number;
ntime: string | number;
version: string | number;
extraNonce2: string;
isBlockCandidate: boolean;
blockSubmissionResult?: string | null;
}
export interface ShareAccountingSummary {
totalAcceptedShares: number;
totalCreditedDifficulty: number;
acceptedSharesLast10Minutes: number;
creditedDifficultyLast10Minutes: number;
acceptedSharesLastHour: number;
creditedDifficultyLastHour: number;
acceptedSharesLastDay: number;
creditedDifficultyLastDay: number;
hashRateLast10Minutes: number;
hashRateLastHour: number;
bestSubmissionDifficulty: number;
blockCandidateCount: number;
latestShareAt: string | null;
protocolBreakdown: {
protocol: string;
acceptedShares: number;
creditedDifficulty: number;
}[];
}
export interface SessionShareSummary {
clientId: string;
latestShareAt: string | null;
hashRateLast10Minutes: number;
bestSubmissionDifficulty: number;
}
interface AccountingFilter {
address?: string;
clientName?: string;
clientId?: string;
}
const HASHES_PER_DIFFICULTY = 4294967296;
const DEFAULT_BATCH_SIZE = 500;
const DEFAULT_FLUSH_INTERVAL_MS = 25;
const DEFAULT_MAX_QUEUE_SIZE = 50000;
const DEFAULT_SUMMARY_CACHE_MS = 2500;
const DEFAULT_SUMMARY_CACHE_MAX = 10000;
@Injectable()
export class ShareAccountingService implements OnModuleDestroy {
private pendingShares: PendingShare[] = [];
private flushTimer: NodeJS.Timeout | null = null;
private activeFlush: Promise<void> | null = null;
private summaryCache = new Map<string, SummaryCacheEntry>();
private readonly batchSize = this.readPositiveInt('SHARE_ACCOUNTING_BATCH_SIZE', DEFAULT_BATCH_SIZE);
private readonly flushIntervalMs = this.readPositiveInt('SHARE_ACCOUNTING_FLUSH_INTERVAL_MS', DEFAULT_FLUSH_INTERVAL_MS);
private readonly maxQueueSize = this.readPositiveInt('SHARE_ACCOUNTING_MAX_QUEUE_SIZE', DEFAULT_MAX_QUEUE_SIZE);
private readonly summaryCacheMs = this.readNonNegativeInt('SHARE_ACCOUNTING_SUMMARY_CACHE_MS', DEFAULT_SUMMARY_CACHE_MS);
private readonly summaryCacheMax = this.readPositiveInt('SHARE_ACCOUNTING_SUMMARY_CACHE_MAX', DEFAULT_SUMMARY_CACHE_MAX);
constructor(
@InjectRepository(AcceptedShareEntity)
private readonly acceptedShareRepository: Repository<AcceptedShareEntity>,
) { }
public async recordAcceptedShare(record: AcceptedShareRecord): Promise<AcceptedShareEntity> {
const acceptedShare = this.acceptedShareRepository.create({
...record,
acceptedAt: record.acceptedAt ?? new Date(),
nonce: record.nonce.toString(),
ntime: record.ntime.toString(),
version: record.version.toString(),
blockSubmissionResult: record.blockSubmissionResult == null
? null
: record.blockSubmissionResult.toString(),
});
if (this.pendingShares.length >= this.maxQueueSize) {
throw new Error(`Share accounting queue is full (${this.maxQueueSize})`);
}
return await new Promise<AcceptedShareEntity>((resolve, reject) => {
this.pendingShares.push({ entity: acceptedShare, resolve, reject });
if (this.pendingShares.length >= this.batchSize) {
this.startFlush();
return;
}
this.scheduleFlush();
});
}
public async flushPendingShares(): Promise<void> {
while (this.pendingShares.length > 0 || this.activeFlush != null) {
if (this.activeFlush != null) {
await this.activeFlush;
continue;
}
this.startFlush();
}
}
public async onModuleDestroy(): Promise<void> {
await this.flushPendingShares();
}
public async getPoolSummary(): Promise<ShareAccountingSummary> {
return this.getSummary({});
}
public async getAddressSummary(address: string): Promise<ShareAccountingSummary> {
return this.getSummary({ address });
}
public async getWorkerGroupSummary(address: string, clientName: string): Promise<ShareAccountingSummary> {
return this.getSummary({ address, clientName });
}
public async getSessionSummary(clientId: string): Promise<ShareAccountingSummary> {
return this.getSummary({ clientId });
}
public async getSessionSummaries(clientIds: string[]): Promise<Map<string, SessionShareSummary>> {
const uniqueClientIds = [...new Set(clientIds.filter(clientId => clientId != null))];
const summaries = new Map<string, SessionShareSummary>();
if (uniqueClientIds.length === 0) {
return summaries;
}
const rows = await this.acceptedShareRepository.query(`
SELECT
"clientId",
MAX("acceptedAt") AS "latestShareAt",
COALESCE((SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes') * ${HASHES_PER_DIFFICULTY}) / 600, 0)::float AS "hashRateLast10Minutes",
COALESCE(MAX("submissionDifficulty"), 0)::float AS "bestSubmissionDifficulty"
FROM "accepted_share_entity"
WHERE "clientId" = ANY($1::uuid[])
GROUP BY "clientId"
`, [uniqueClientIds]);
rows.forEach(row => {
summaries.set(row.clientId, {
clientId: row.clientId,
latestShareAt: row.latestShareAt == null
? null
: new Date(row.latestShareAt).toISOString(),
hashRateLast10Minutes: this.toNumber(row.hashRateLast10Minutes),
bestSubmissionDifficulty: this.toNumber(row.bestSubmissionDifficulty),
});
});
return summaries;
}
private async getSummary(filter: AccountingFilter): Promise<ShareAccountingSummary> {
const cacheKey = this.getSummaryCacheKey(filter);
const cached = this.summaryCache.get(cacheKey);
const now = Date.now();
if (cached != null && cached.expiresAt > now) {
return cached.value;
}
const value = this.loadSummary(filter).catch(error => {
this.summaryCache.delete(cacheKey);
throw error;
});
if (this.summaryCacheMs > 0) {
this.summaryCache.set(cacheKey, {
expiresAt: now + this.summaryCacheMs,
value,
});
this.trimSummaryCache();
}
return value;
}
private async loadSummary(filter: AccountingFilter): Promise<ShareAccountingSummary> {
const { whereSql, params } = this.buildWhereClause(filter);
const [summary] = await this.acceptedShareRepository.query(`
SELECT
COUNT(*)::int AS "totalAcceptedShares",
COALESCE(SUM("creditedDifficulty"), 0)::float AS "totalCreditedDifficulty",
COUNT(*) FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes')::int AS "acceptedSharesLast10Minutes",
COALESCE(SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes'), 0)::float AS "creditedDifficultyLast10Minutes",
COUNT(*) FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 hour')::int AS "acceptedSharesLastHour",
COALESCE(SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 hour'), 0)::float AS "creditedDifficultyLastHour",
COUNT(*) FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 day')::int AS "acceptedSharesLastDay",
COALESCE(SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 day'), 0)::float AS "creditedDifficultyLastDay",
COALESCE((SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes') * ${HASHES_PER_DIFFICULTY}) / 600, 0)::float AS "hashRateLast10Minutes",
COALESCE((SUM("creditedDifficulty") FILTER (WHERE "acceptedAt" > NOW() - INTERVAL '1 hour') * ${HASHES_PER_DIFFICULTY}) / 3600, 0)::float AS "hashRateLastHour",
COALESCE(MAX("submissionDifficulty"), 0)::float AS "bestSubmissionDifficulty",
COUNT(*) FILTER (WHERE "isBlockCandidate" = TRUE)::int AS "blockCandidateCount",
MAX("acceptedAt") AS "latestShareAt"
FROM "accepted_share_entity"
${whereSql}
`, params);
const protocolRows = await this.acceptedShareRepository.query(`
SELECT
"protocol",
COUNT(*)::int AS "acceptedShares",
COALESCE(SUM("creditedDifficulty"), 0)::float AS "creditedDifficulty"
FROM "accepted_share_entity"
${whereSql}
GROUP BY "protocol"
ORDER BY "protocol"
`, params);
return {
totalAcceptedShares: this.toNumber(summary?.totalAcceptedShares),
totalCreditedDifficulty: this.toNumber(summary?.totalCreditedDifficulty),
acceptedSharesLast10Minutes: this.toNumber(summary?.acceptedSharesLast10Minutes),
creditedDifficultyLast10Minutes: this.toNumber(summary?.creditedDifficultyLast10Minutes),
acceptedSharesLastHour: this.toNumber(summary?.acceptedSharesLastHour),
creditedDifficultyLastHour: this.toNumber(summary?.creditedDifficultyLastHour),
acceptedSharesLastDay: this.toNumber(summary?.acceptedSharesLastDay),
creditedDifficultyLastDay: this.toNumber(summary?.creditedDifficultyLastDay),
hashRateLast10Minutes: this.toNumber(summary?.hashRateLast10Minutes),
hashRateLastHour: this.toNumber(summary?.hashRateLastHour),
bestSubmissionDifficulty: this.toNumber(summary?.bestSubmissionDifficulty),
blockCandidateCount: this.toNumber(summary?.blockCandidateCount),
latestShareAt: summary?.latestShareAt == null
? null
: new Date(summary.latestShareAt).toISOString(),
protocolBreakdown: protocolRows.map(row => ({
protocol: row.protocol,
acceptedShares: this.toNumber(row.acceptedShares),
creditedDifficulty: this.toNumber(row.creditedDifficulty),
})),
};
}
private scheduleFlush(): void {
if (this.flushTimer != null || this.activeFlush != null) {
return;
}
this.flushTimer = setTimeout(() => {
this.flushTimer = null;
this.startFlush();
}, this.flushIntervalMs);
this.flushTimer.unref?.();
}
private startFlush(): void {
if (this.activeFlush != null) {
return;
}
if (this.flushTimer != null) {
clearTimeout(this.flushTimer);
this.flushTimer = null;
}
const batch = this.pendingShares.splice(0, this.batchSize);
if (batch.length === 0) {
return;
}
this.activeFlush = this.flushBatch(batch).finally(() => {
this.activeFlush = null;
if (this.pendingShares.length >= this.batchSize) {
this.startFlush();
return;
}
if (this.pendingShares.length > 0) {
this.scheduleFlush();
}
});
}
private async flushBatch(batch: PendingShare[]): Promise<void> {
try {
await this.acceptedShareRepository.insert(batch.map(item => item.entity));
batch.forEach(item => item.resolve(item.entity));
} catch (error) {
batch.forEach(item => item.reject(error));
}
}
private buildWhereClause(filter: AccountingFilter): { whereSql: string; params: string[] } {
const where: string[] = [];
const params: string[] = [];
if (filter.address != null) {
params.push(filter.address);
where.push(`"address" = $${params.length}`);
}
if (filter.clientName != null) {
params.push(filter.clientName);
where.push(`"clientName" = $${params.length}`);
}
if (filter.clientId != null) {
params.push(filter.clientId);
where.push(`"clientId" = $${params.length}`);
}
return {
whereSql: where.length > 0 ? `WHERE ${where.join(' AND ')}` : '',
params,
};
}
private toNumber(value: unknown): number {
const parsed = Number(value ?? 0);
return Number.isFinite(parsed) ? parsed : 0;
}
private getSummaryCacheKey(filter: AccountingFilter): string {
return JSON.stringify({
address: filter.address ?? null,
clientName: filter.clientName ?? null,
clientId: filter.clientId ?? null,
});
}
private trimSummaryCache(): void {
while (this.summaryCache.size > this.summaryCacheMax) {
const firstKey = this.summaryCache.keys().next().value;
if (firstKey == null) {
return;
}
this.summaryCache.delete(firstKey);
}
}
private readPositiveInt(name: string, defaultValue: number): number {
const value = Number(process.env[name]);
return Number.isInteger(value) && value > 0 ? value : defaultValue;
}
private readNonNegativeInt(name: string, defaultValue: number): number {
const value = Number(process.env[name]);
return Number.isInteger(value) && value >= 0 ? value : defaultValue;
}
}
interface PendingShare {
entity: AcceptedShareEntity;
resolve: (entity: AcceptedShareEntity) => void;
reject: (error: unknown) => void;
}
interface SummaryCacheEntry {
expiresAt: number;
value: Promise<ShareAccountingSummary>;
}
+29 -9
View File
@@ -8,10 +8,10 @@ import { AddressSettingsService } from './ORM/address-settings/address-settings.
import { BlocksService } from './ORM/blocks/blocks.service';
import { ClientStatisticsService } from './ORM/client-statistics/client-statistics.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 { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-report.view';
import { StratumV2Service } from './services/stratum-v2.service';
import { ShareAccountingService } from './ORM/share-accounting/share-accounting.service';
@Controller()
export class AppController {
@@ -24,10 +24,10 @@ export class AppController {
private readonly clientStatisticsService: ClientStatisticsService,
private readonly blocksService: BlocksService,
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly homeGraphService: HomeGraphService,
private readonly addressSettingsService: AddressSettingsService,
private readonly userAgentReportService: UserAgentReportService,
private readonly stratumV2Service: StratumV2Service
private readonly stratumV2Service: StratumV2Service,
private readonly shareAccountingService: ShareAccountingService
) { }
@Get('info')
@@ -69,12 +69,15 @@ export class AppController {
return pre;
}, []);
userAgents.push({ userAgent: 'Other', count: other.count.toString(), bestDifficulty: other.bestDifficulty, totalHashRate: other.totalHashRate.toString() })
if (other.count > 0) {
userAgents.push({ userAgent: 'Other', count: other.count.toString(), bestDifficulty: other.bestDifficulty, totalHashRate: other.totalHashRate.toString() })
}
const data = {
blockData,
userAgents,
highScores,
accounting: await this.shareAccountingService.getPoolSummary(),
sv2: {
poolAuthorityPublicKey: poolAuthority.publicKey,
authorityKeyConfigured: poolAuthority.configured
@@ -82,13 +85,30 @@ export class AppController {
uptime: this.uptime
};
//5 min
await this.cacheManager.set(CACHE_KEY, data, 5 * 60 * 1000);
// Keep online miner counts responsive after reconnect cleanup.
await this.cacheManager.set(CACHE_KEY, data, 15 * 1000);
return data;
}
@Get('info/accounting')
public async infoAccounting() {
const CACHE_KEY = 'SITE_ACCOUNTING';
const cachedResult = await this.cacheManager.get(CACHE_KEY);
if (cachedResult != null) {
return cachedResult;
}
const data = await this.shareAccountingService.getPoolSummary();
//15 sec
await this.cacheManager.set(CACHE_KEY, data, 15 * 1000);
return data;
}
@Get('pool')
public async pool() {
@@ -114,8 +134,8 @@ export class AppController {
fee: 0
}
//5 min
await this.cacheManager.set(CACHE_KEY, data, 5 * 60 * 1000);
// Keep online miner counts responsive after reconnect cleanup.
await this.cacheManager.set(CACHE_KEY, data, 15 * 1000);
return data;
}
@@ -136,7 +156,7 @@ export class AppController {
return cachedResult;
}
const chartData = await this.homeGraphService.getChartDataForSite();
const chartData = await this.clientStatisticsService.getChartDataForSite();
//10 min
await this.cacheManager.set(CACHE_KEY, chartData, 10 * 60 * 1000);
+16 -36
View File
@@ -9,29 +9,22 @@ import { AppController } from './app.controller';
import { AddressController } from './controllers/address/address.controller';
import { ClientController } from './controllers/client/client.controller';
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 { BlocksEntity } from './ORM/blocks/blocks.entity';
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 { ClientEntity } from './ORM/client/client.entity';
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 { TelegramSubscriptionsEntity } from './ORM/telegram-subscriptions/telegram-subscriptions.entity';
import { ShareAccountingModule } from './ORM/share-accounting/share-accounting.module';
import { TelegramSubscriptionsModule } from './ORM/telegram-subscriptions/telegram-subscriptions.module';
import { createDatabaseOptions } from './database.config';
import { AppService } from './services/app.service';
import { BitcoinRpcService } from './services/bitcoin-rpc.service';
import { BraiinsService } from './services/braiins.service';
import { BTCPayService } from './services/btc-pay.service';
import { DiscordService } from './services/discord.service';
import { NotificationService } from './services/notification.service';
import { RedisMessagingModule } from './services/redis-messaging.module';
import { StratumV1JobsService } from './services/stratum-v1-jobs.service';
import { StratumV1Service } from './services/stratum-v1.service';
import { StratumV2Service } from './services/stratum-v2.service';
@@ -45,8 +38,8 @@ const ORMModules = [
TelegramSubscriptionsModule,
BlocksModule,
RpcBlocksModule,
HomeGraphModule,
UserAgentReportModule
UserAgentReportModule,
ShareAccountingModule
]
@Module({
@@ -54,30 +47,16 @@ const ORMModules = [
ConfigModule.forRoot(),
TypeOrmModule.forRootAsync({
useFactory: (configService: ConfigService) => {
return {
type: 'postgres',
host: configService.get('DB_HOST'),
port: parseInt(configService.get('DB_PORT')),
username: configService.get('DB_USERNAME'),
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
]
}
return createDatabaseOptions({
...process.env,
DB_HOST: configService.get('DB_HOST'),
DB_PORT: configService.get('DB_PORT'),
DB_USERNAME: configService.get('DB_USERNAME'),
DB_PASSWORD: configService.get('DB_PASSWORD'),
DB_DATABASE: configService.get('DB_DATABASE'),
DB_LOGGING: configService.get('DB_LOGGING'),
DB_POOL_SIZE: configService.get('DB_POOL_SIZE'),
});
},
imports: [ConfigModule],
inject: [ConfigService]
@@ -85,6 +64,7 @@ const ORMModules = [
CacheModule.register(),
ScheduleModule.forRoot(),
HttpModule,
RedisMessagingModule,
...ORMModules
],
controllers: [
@@ -1,32 +1,59 @@
import { Test, TestingModule } from '@nestjs/testing';
import { TypeOrmModule } from '@nestjs/typeorm';
import { AddressSettingsModule } from '../../ORM/address-settings/address-settings.module';
import { ClientStatisticsModule } from '../../ORM/client-statistics/client-statistics.module';
import { ClientModule } from '../../ORM/client/client.module';
import { AddressSettingsService } from '../../ORM/address-settings/address-settings.service';
import { ClientStatisticsService } from '../../ORM/client-statistics/client-statistics.service';
import { ClientService } from '../../ORM/client/client.service';
import { ShareAccountingService } from '../../ORM/share-accounting/share-accounting.service';
import { RedisMessagingService } from '../../services/redis-messaging.service';
import { ClientController } from './client.controller';
describe('ClientController', () => {
let controller: ClientController;
beforeEach(async () => {
const module: TestingModule = await Test.createTestingModule({
imports: [
TypeOrmModule.forRoot({
type: 'better-sqlite3',
database: ':memory:',
synchronize: true,
autoLoadEntities: true,
cache: true,
logging: false
}),
AddressSettingsModule,
ClientModule,
ClientStatisticsModule
],
controllers: [ClientController],
beforeEach(async () => {
const module: TestingModule = await Test.createTestingModule({
controllers: [ClientController],
providers: [
{
provide: ClientService,
useValue: {
getByAddress: jest.fn(),
getByName: jest.fn(),
getBySessionId: jest.fn(),
},
},
{
provide: ClientStatisticsService,
useValue: {
getChartDataForAddress: jest.fn(),
getChartDataForGroup: jest.fn(),
getChartDataForSession: jest.fn(),
},
},
{
provide: AddressSettingsService,
useValue: {
getSettings: jest.fn(),
},
},
{
provide: ShareAccountingService,
useValue: {
getAddressSummary: jest.fn(),
getWorkerGroupSummary: jest.fn(),
getSessionSummary: jest.fn(),
getSessionSummaries: jest.fn().mockResolvedValue(new Map()),
},
},
{
provide: RedisMessagingService,
useValue: {
getClientPresenceByAddress: jest.fn().mockResolvedValue([]),
},
},
],
}).compile();
}).compile();
controller = module.get<ClientController>(ClientController);
});
+31 -7
View File
@@ -3,6 +3,8 @@ import { Controller, Get, NotFoundException, Param } from '@nestjs/common';
import { AddressSettingsService } from '../../ORM/address-settings/address-settings.service';
import { ClientStatisticsService } from '../../ORM/client-statistics/client-statistics.service';
import { ClientService } from '../../ORM/client/client.service';
import { ShareAccountingService } from '../../ORM/share-accounting/share-accounting.service';
import { RedisMessagingService } from '../../services/redis-messaging.service';
@Controller('client')
@@ -11,29 +13,38 @@ export class ClientController {
constructor(
private readonly clientService: ClientService,
private readonly clientStatisticsService: ClientStatisticsService,
private readonly addressSettingsService: AddressSettingsService
private readonly addressSettingsService: AddressSettingsService,
private readonly shareAccountingService: ShareAccountingService,
private readonly redisMessagingService: RedisMessagingService
) { }
@Get(':address')
async getClientInfo(@Param('address') address: string) {
const workers = await this.clientService.getByAddress(address);
const workers = await this.redisMessagingService.getClientPresenceByAddress(address);
const sessionSummaries = await this.shareAccountingService.getSessionSummaries(workers.map(worker => worker.clientId));
const addressSettings = await this.addressSettingsService.getSettings(address, false);
return {
bestDifficulty: addressSettings?.bestDifficulty,
workersCount: workers.length,
accounting: await this.shareAccountingService.getAddressSummary(address),
workers: await Promise.all(
workers.map(async (worker) => {
const sessionSummary = sessionSummaries.get(worker.clientId);
const bestDifficulty = Math.max(
Number(worker.bestDifficulty ?? 0),
Number(sessionSummary?.bestSubmissionDifficulty ?? 0),
);
return {
sessionId: worker.sessionId,
name: worker.clientName,
bestDifficulty: parseFloat(worker.bestDifficulty as any).toFixed(2),
hashRate: worker.hashRate,
bestDifficulty: bestDifficulty.toFixed(2),
hashRate: sessionSummary?.hashRateLast10Minutes ?? worker.hashRate,
startTime: worker.startTime,
lastSeen: worker.updatedAt
lastSeen: sessionSummary?.latestShareAt ?? worker.lastSeen
};
})
)
@@ -49,7 +60,8 @@ export class ClientController {
@Get(':address/:workerName')
async getWorkerGroupInfo(@Param('address') address: string, @Param('workerName') workerName: string) {
const workers = await this.clientService.getByName(address, workerName);
const workers = (await this.redisMessagingService.getClientPresenceByAddress(address))
.filter(worker => worker.clientName === workerName);
const bestDifficulty = workers.reduce((pre, cur, idx, arr) => {
if (cur.bestDifficulty > pre) {
@@ -63,6 +75,7 @@ export class ClientController {
name: workerName,
bestDifficulty: Math.floor(bestDifficulty),
accounting: await this.shareAccountingService.getWorkerGroupSummary(address, workerName),
chartData: chartData,
}
@@ -71,7 +84,17 @@ export class ClientController {
@Get(':address/:workerName/:sessionId')
async getWorkerInfo(@Param('address') address: string, @Param('workerName') workerName: string, @Param('sessionId') sessionId: string) {
const worker = await this.clientService.getBySessionId(address, workerName, sessionId);
const presenceWorker = (await this.redisMessagingService.getClientPresenceByAddress(address))
.find(worker => worker.clientName === workerName && worker.sessionId === sessionId);
const worker = presenceWorker == null
? await this.clientService.getBySessionId(address, workerName, sessionId)
: {
id: presenceWorker.clientId,
sessionId: presenceWorker.sessionId,
clientName: presenceWorker.clientName,
bestDifficulty: presenceWorker.bestDifficulty,
startTime: presenceWorker.startTime,
};
if (worker == null) {
return new NotFoundException();
}
@@ -81,6 +104,7 @@ export class ClientController {
sessionId: worker.sessionId,
name: worker.clientName,
bestDifficulty: Math.floor(worker.bestDifficulty),
accounting: await this.shareAccountingService.getSessionSummary(worker.id),
chartData: chartData,
startTime: worker.startTime
}
+7
View File
@@ -0,0 +1,7 @@
import 'dotenv/config';
import 'reflect-metadata';
import { DataSource } from 'typeorm';
import { createDatabaseOptions } from './database.config';
export const AppDataSource = new DataSource(createDatabaseOptions(process.env));
+49
View File
@@ -0,0 +1,49 @@
import { TypeOrmModuleOptions } from '@nestjs/typeorm';
import { DataSourceOptions } from 'typeorm';
import { InitialTimescaleSchema1780859300000 } from './ORM/_migrations/InitialTimescaleSchema1780859300000';
import { ActiveOnlyUserAgentReport1780860200000 } from './ORM/_migrations/ActiveOnlyUserAgentReport1780860200000';
import { TimescaleOperationalHardening1780861200000 } from './ORM/_migrations/TimescaleOperationalHardening1780861200000';
import { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-report.view';
import { AcceptedShareEntity } from './ORM/accepted-share/accepted-share.entity';
import { AddressSettingsEntity } from './ORM/address-settings/address-settings.entity';
import { BlocksEntity } from './ORM/blocks/blocks.entity';
import { ClientEntity } from './ORM/client/client.entity';
import { RpcBlockEntity } from './ORM/rpc-block/rpc-block.entity';
import { TelegramSubscriptionsEntity } from './ORM/telegram-subscriptions/telegram-subscriptions.entity';
export const databaseEntities = [
ClientEntity,
AddressSettingsEntity,
BlocksEntity,
RpcBlockEntity,
TelegramSubscriptionsEntity,
UserAgentReportView,
AcceptedShareEntity,
];
export const databaseMigrations = [
InitialTimescaleSchema1780859300000,
ActiveOnlyUserAgentReport1780860200000,
TimescaleOperationalHardening1780861200000,
];
export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions {
return {
type: 'postgres',
host: env.DB_HOST,
port: parseInt(env.DB_PORT ?? '5432', 10),
username: env.DB_USERNAME,
password: env.DB_PASSWORD,
database: env.DB_DATABASE,
entities: databaseEntities,
synchronize: false,
logging: env.DB_LOGGING === 'true',
poolSize: parseInt(env.DB_POOL_SIZE ?? '10', 10),
ssl: env.DB_SSL === 'true'
? { rejectUnauthorized: env.DB_SSL_REJECT_UNAUTHORIZED !== 'false' }
: undefined,
migrations: databaseMigrations,
migrationsTransactionMode: 'none',
};
}
+66 -80
View File
@@ -1,20 +1,10 @@
import { ConfigService } from '@nestjs/config';
import { Test, TestingModule } from '@nestjs/testing';
import { TypeOrmModule } from '@nestjs/typeorm';
import { Socket } from 'net';
import { BehaviorSubject } from 'rxjs';
import { DataSource } from 'typeorm';
import { MockRecording1 } from '../../test/models/MockRecording1';
import { AddressSettingsModule } from '../ORM/address-settings/address-settings.module';
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
import { BlocksEntity } from '../ORM/blocks/blocks.entity';
import { BlocksService } from '../ORM/blocks/blocks.service';
import { ClientStatisticsEntity } from '../ORM/client-statistics/client-statistics.entity';
import { ClientStatisticsModule } from '../ORM/client-statistics/client-statistics.module';
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
import { ClientEntity } from '../ORM/client/client.entity';
import { ClientModule } from '../ORM/client/client.module';
import { ClientService } from '../ORM/client/client.service';
import { BitcoinRpcService as MockBitcoinRpcService } from '../services/bitcoin-rpc.service';
import { NotificationService } from '../services/notification.service';
@@ -45,10 +35,12 @@ describe('StratumV1Client', () => {
let bitcoinRpcService: MockBitcoinRpcService;
let clientService: ClientService;
let clientStatisticsService: ClientStatisticsService;
let notificationService: NotificationService;
let blocksService: BlocksService;
let configService: ConfigService;
let addressSettings: AddressSettingsService;
let shareAccountingService: { recordAcceptedShare: jest.Mock };
let redisMessagingService: { setClientPresence: jest.Mock; removeClientPresence: jest.Mock };
let client: StratumV1Client;
@@ -60,46 +52,6 @@ describe('StratumV1Client', () => {
let newBlockEmitter: BehaviorSubject<IBlockTemplate> = new BehaviorSubject(MockRecording1.BLOCK_TEMPLATE);
let moduleRef: TestingModule;
beforeAll(async () => {
moduleRef = await Test.createTestingModule({
imports: [
TypeOrmModule.forRoot({
type: 'better-sqlite3',
database: ':memory:',
synchronize: true,
autoLoadEntities: true,
cache: true,
logging: false,
entities: [ClientEntity, ClientStatisticsEntity, BlocksEntity]
}),
ClientModule,
ClientStatisticsModule,
AddressSettingsModule
],
providers: [
{
provide: ConfigService,
useValue: {
get: jest.fn((key: string) => {
switch (key) {
case 'DEV_FEE_ADDRESS':
return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
case 'NETWORK':
return 'testnet';
}
return null;
})
}
}
],
}).compile();
})
beforeEach(async () => {
jest.useFakeTimers({ advanceTimers: true })
@@ -108,27 +60,38 @@ describe('StratumV1Client', () => {
consoleErrorSpy = jest.spyOn(console, 'error').mockImplementation(() => undefined);
consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined);
clientService = moduleRef.get<ClientService>(ClientService);
const clients = new Map<string, any>();
let nextClientId = 1;
clientService = {
insert: jest.fn(async (partialClient) => {
const clientEntity = {
id: `00000000-0000-4000-8000-${String(nextClientId++).padStart(12, '0')}`,
hashRate: 0,
updatedAt: new Date(),
deletedAt: null,
...partialClient,
};
clients.set(clientEntity.id, clientEntity);
return clientEntity;
}),
delete: jest.fn(async (id: string) => {
clients.delete(id);
}),
connectedClientCount: jest.fn(async () => clients.size),
updateBestDifficultyIfHigher: jest.fn().mockResolvedValue({ affected: 1 }),
} as any;
const dataSource = moduleRef.get<DataSource>(DataSource);
await dataSource.getRepository(ClientStatisticsEntity).clear();
await dataSource.getRepository(ClientEntity).clear();
await dataSource.getRepository(BlocksEntity).clear();
clientStatisticsService = moduleRef.get<ClientStatisticsService>(ClientStatisticsService);
configService = moduleRef.get<ConfigService>(ConfigService);
(configService.get as jest.Mock).mockImplementation((key: string) => {
switch (key) {
case 'DEV_FEE_ADDRESS':
return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
case 'NETWORK':
return 'testnet';
}
return null;
});
configService = {
get: jest.fn((key: string) => {
switch (key) {
case 'DEV_FEE_ADDRESS':
return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
case 'NETWORK':
return 'testnet';
}
return null;
})
} as any;
(StratumV1Client as any).blockedUserAgentLogState.clear();
(StratumV1Client as any).validationErrorLogState.clear();
@@ -154,13 +117,24 @@ describe('StratumV1Client', () => {
socket.end = jest.fn();
jest.spyOn(socket, 'destroy').mockImplementation(() => socket);
const addressSettings = moduleRef.get<AddressSettingsService>(AddressSettingsService);
addressSettings = {
getSettings: jest.fn().mockResolvedValue(null),
updateBestDifficultyIfHigher: jest.fn().mockResolvedValue({ affected: 1 }),
resetBestDifficultyAndShares: jest.fn().mockResolvedValue(undefined),
} as any;
notificationService = {
notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined)
} as any;
blocksService = {
save: jest.fn().mockResolvedValue(undefined)
} as any;
shareAccountingService = {
recordAcceptedShare: jest.fn().mockResolvedValue(undefined),
};
redisMessagingService = {
setClientPresence: jest.fn().mockResolvedValue(undefined),
removeClientPresence: jest.fn().mockResolvedValue(undefined),
};
client = new StratumV1Client(
@@ -168,11 +142,12 @@ describe('StratumV1Client', () => {
stratumV1JobsService,
bitcoinRpcService,
clientService,
clientStatisticsService,
notificationService,
blocksService,
configService,
addressSettings
addressSettings,
shareAccountingService as any,
redisMessagingService as any
);
client.extraNonceAndSessionId = MockRecording1.EXTRA_NONCE;
@@ -263,11 +238,10 @@ describe('StratumV1Client', () => {
stratumV1JobsService,
bitcoinRpcService,
clientService,
clientStatisticsService,
notificationService,
blocksService,
configService,
moduleRef.get<AddressSettingsService>(AddressSettingsService)
addressSettings
);
socketEmitter(Buffer.from(`{"id":1,"method":"mining.subscribe","params":["NMMiner/1.0"]}\n`));
@@ -369,6 +343,21 @@ describe('StratumV1Client', () => {
await new Promise((r) => setTimeout(r, 1000));
expect((client as any).write).lastCalledWith(`{\"id\":5,\"error\":null,\"result\":true}\n`);
expect(shareAccountingService.recordAcceptedShare).toHaveBeenCalledWith(expect.objectContaining({
protocol: 'sv1',
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
clientName: 'bitaxe3',
sessionId: MockRecording1.EXTRA_NONCE,
jobId: '1',
jobTemplateId: '1',
creditedDifficulty: 0,
isBlockCandidate: false,
}));
expect(redisMessagingService.setClientPresence).toHaveBeenCalledWith(expect.objectContaining({
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
clientName: 'bitaxe3',
sessionId: MockRecording1.EXTRA_NONCE,
}));
});
@@ -423,7 +412,6 @@ describe('StratumV1Client', () => {
submissionDifficulty: 1024,
submissionHash: 'share'
});
const addressSettings = moduleRef.get<AddressSettingsService>(AddressSettingsService);
const getSettingsSpy = jest.spyOn(addressSettings, 'getSettings');
const updateIfHigherSpy = jest.spyOn(addressSettings as any, 'updateBestDifficultyIfHigher').mockResolvedValue({ affected: 1 });
const clientUpdateIfHigherSpy = jest.spyOn(clientService as any, 'updateBestDifficultyIfHigher').mockResolvedValue({ affected: 1 });
@@ -546,11 +534,10 @@ describe('StratumV1Client', () => {
stratumV1JobsService,
bitcoinRpcService,
clientService,
clientStatisticsService,
notificationService,
blocksService,
configService,
moduleRef.get<AddressSettingsService>(AddressSettingsService)
addressSettings
);
jest.spyOn(secondClient as any, 'write').mockImplementation((data) => Promise.resolve(true));
jest.spyOn(secondClient as any, 'getRandomHexString').mockReturnValue(MockRecording1.EXTRA_NONCE);
@@ -580,7 +567,6 @@ describe('StratumV1Client', () => {
submissionDifficulty: Number.MAX_SAFE_INTEGER,
submissionHash: 'block-share'
});
const addressSettings = moduleRef.get<AddressSettingsService>(AddressSettingsService);
jest.spyOn(addressSettings, 'resetBestDifficultyAndShares').mockResolvedValue(undefined);
emitMessage(MockRecording1.MINING_SUBSCRIBE);
+58 -12
View File
@@ -9,11 +9,12 @@ import { clearInterval } from 'timers';
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
import { BlocksService } from '../ORM/blocks/blocks.service';
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
import { ClientEntity } from '../ORM/client/client.entity';
import { ClientService } from '../ORM/client/client.service';
import { ShareAccountingService } from '../ORM/share-accounting/share-accounting.service';
import { BitcoinRpcService } from '../services/bitcoin-rpc.service';
import { NotificationService } from '../services/notification.service';
import { RedisMessagingService } from '../services/redis-messaging.service';
import { IJobTemplate, StratumV1JobsService } from '../services/stratum-v1-jobs.service';
import { eRequestMethod } from './enums/eRequestMethod';
import { eResponseMethod } from './enums/eResponseMethod';
@@ -67,11 +68,12 @@ export class StratumV1Client {
private readonly stratumV1JobsService: StratumV1JobsService,
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly clientService: ClientService,
private readonly clientStatisticsService: ClientStatisticsService,
private readonly notificationService: NotificationService,
private readonly blocksService: BlocksService,
private readonly configService: ConfigService,
private readonly addressSettingsService: AddressSettingsService
private readonly addressSettingsService: AddressSettingsService,
private readonly shareAccountingService?: ShareAccountingService,
private readonly redisMessagingService?: RedisMessagingService
) {
this.socket.on('data', (data: Buffer) => {
@@ -100,6 +102,7 @@ export class StratumV1Client {
public async destroy() {
if (this.clientEntity?.id) {
await this.redisMessagingService?.removeClientPresence(this.clientEntity.id, this.clientEntity.address);
await this.clientService.delete(this.clientEntity.id);
}
@@ -156,7 +159,7 @@ export class StratumV1Client {
if (this.sessionStart == null) {
this.sessionStart = new Date();
this.statistics = new StratumV1ClientStatistics(this.clientStatisticsService);
this.statistics = new StratumV1ClientStatistics();
this.extraNonceAndSessionId = this.getRandomHexString();
//console.log(`New client ID: : ${this.extraNonceAndSessionId}, ${this.socket.remoteAddress}:${this.socket.remotePort}`);
}
@@ -422,7 +425,6 @@ export class StratumV1Client {
// //50Th/s
// this.noFee = false;
// if (this.clientEntity) {
// this.hashRate = await this.clientStatisticsService.getHashRateForSession(this.clientEntity.id);
// // 250Gh/s
// if(this.hashRate < 250000000000){
// this.statistics.targetSubmitShareEveryNSeconds = 10;
@@ -491,12 +493,35 @@ export class StratumV1Client {
startTime: new Date(),
bestDifficulty: 0
});
await this.updateClientPresence(new Date());
})();
}
await this.creatingEntity;
}
private async updateClientPresence(lastSeen: Date): Promise<void> {
if (this.clientEntity == null) {
return;
}
try {
await this.redisMessagingService?.setClientPresence({
clientId: this.clientEntity.id,
address: this.clientEntity.address,
clientName: this.clientEntity.clientName,
sessionId: this.clientEntity.sessionId,
userAgent: this.clientEntity.userAgent,
startTime: new Date(this.clientEntity.startTime).toISOString(),
lastSeen: lastSeen.toISOString(),
hashRate: this.statistics?.hashRate ?? 0,
bestDifficulty: Number(this.clientEntity.bestDifficulty ?? 0),
});
} catch (error) {
console.error(`Failed to update SV1 client presence: ${error.message}`);
}
}
private async handleMiningSubmission(submission: MiningSubmitMessage) {
const job = this.stratumV1JobsService.getJobById(submission.jobId);
@@ -574,6 +599,7 @@ export class StratumV1Client {
return false;
}
let blockSubmissionResult: string = null;
if (submissionDifficulty >= jobTemplate.blockData.networkDifficulty) {
console.log('!!! BLOCK FOUND !!!');
const updatedJobBlock = job.copyAndUpdateBlock(
@@ -585,7 +611,7 @@ export class StratumV1Client {
timestamp
);
const blockHex = updatedJobBlock.toHex(false);
const result = await this.bitcoinRpcService.SUBMIT_BLOCK(blockHex);
blockSubmissionResult = await this.bitcoinRpcService.SUBMIT_BLOCK(blockHex);
await this.blocksService.save({
height: jobTemplate.blockData.height,
minerAddress: this.clientAuthorization.address,
@@ -594,21 +620,40 @@ export class StratumV1Client {
blockData: blockHex
});
await this.notificationService.notifySubscribersBlockFound(this.clientAuthorization.address, jobTemplate.blockData.height, updatedJobBlock, result);
await this.notificationService.notifySubscribersBlockFound(this.clientAuthorization.address, jobTemplate.blockData.height, updatedJobBlock, blockSubmissionResult);
//success
if (result == null) {
if (blockSubmissionResult == null) {
await this.addressSettingsService.resetBestDifficultyAndShares();
}
}
await this.ensureClientEntity();
try {
await this.shareAccountingService?.recordAcceptedShare({
protocol: 'sv1',
address: this.clientAuthorization.address,
clientName: this.clientAuthorization.worker,
sessionId: this.extraNonceAndSessionId,
clientId: this.clientEntity.id,
jobId: job.jobId,
jobTemplateId: job.jobTemplateId,
blockHeight: jobTemplate.blockData.height,
creditedDifficulty: this.sessionDifficulty,
submissionDifficulty,
networkDifficulty: jobTemplate.blockData.networkDifficulty,
nonce: submission.nonce,
ntime: submission.ntime,
version: Number.isFinite(versionMask)
? (jobTemplate.block.version ^ versionMask).toString(16)
: jobTemplate.block.version.toString(16),
extraNonce2: submission.extraNonce2,
isBlockCandidate: submissionDifficulty >= jobTemplate.blockData.networkDifficulty,
blockSubmissionResult,
});
await this.statistics.addShares(this.clientEntity, this.sessionDifficulty);
const now = new Date();
// only update every minute
//if (this.clientEntity.updatedAt == null || now.getTime() - this.clientEntity.updatedAt.getTime() > 1000 * 60) {
this.clientService.heartbeatBulkAsync(this.clientEntity.id, this.statistics.hashRate, now);
this.clientEntity.updatedAt = now;
//}
this.clientEntity.hashRate = this.statistics.hashRate;
await this.updateClientPresence(now);
} catch (e) {
console.log(e);
@@ -618,6 +663,7 @@ export class StratumV1Client {
await this.clientService.updateBestDifficultyIfHigher(this.clientEntity.id, submissionDifficulty);
this.clientEntity.bestDifficulty = submissionDifficulty;
await this.addressSettingsService.updateBestDifficultyIfHigher(this.clientAuthorization.address, submissionDifficulty, this.clientEntity.userAgent);
await this.updateClientPresence(new Date());
}
+9 -51
View File
@@ -8,73 +8,31 @@ describe('StratumV1ClientStatistics', () => {
sessionId: '57a6f098'
} as any;
let clientStatisticsService: {
insert: jest.Mock,
updateBulkAsync: jest.Mock
};
let statistics: StratumV1ClientStatistics;
beforeEach(() => {
jest.useFakeTimers();
jest.setSystemTime(new Date('2026-05-06T12:00:00Z'));
clientStatisticsService = {
insert: jest.fn().mockResolvedValue(undefined),
updateBulkAsync: jest.fn().mockResolvedValue(undefined)
};
statistics = new StratumV1ClientStatistics(clientStatisticsService as any);
statistics = new StratumV1ClientStatistics();
});
afterEach(() => {
jest.useRealTimers();
});
it('should insert the first share bucket', async () => {
it('should keep runtime hashrate at zero before enough time passes', async () => {
await statistics.addShares(client, 64);
expect(clientStatisticsService.insert).toHaveBeenCalledWith({
time: new Date('2026-05-06T12:00:00Z').getTime(),
clientId: client.id,
shares: 64,
acceptedCount: 1,
address: client.address,
clientName: client.clientName,
sessionId: client.sessionId
});
expect(statistics.hashRate).toBe(0);
});
it('should update the current share bucket for additional shares', async () => {
await statistics.addShares(client, 64);
await statistics.addShares(client, 32);
it('should update runtime hashrate from accepted shares', async () => {
for (let i = 0; i < 3; i++) {
jest.setSystemTime(new Date(Date.parse('2026-05-06T12:00:00Z') + (i * 31000)));
await statistics.addShares(client, 64);
}
expect(clientStatisticsService.updateBulkAsync).toHaveBeenCalledWith({
time: new Date('2026-05-06T12:00:00Z').getTime(),
clientId: client.id,
shares: 96,
acceptedCount: 2
});
});
it('should create a new share bucket when the time slot changes', async () => {
await statistics.addShares(client, 64);
jest.setSystemTime(new Date('2026-05-06T12:10:00Z'));
await statistics.addShares(client, 32);
expect(clientStatisticsService.updateBulkAsync).toHaveBeenCalledWith({
time: new Date('2026-05-06T12:00:00Z').getTime(),
clientId: client.id,
shares: 64,
acceptedCount: 1
});
expect(clientStatisticsService.insert).toHaveBeenLastCalledWith({
time: new Date('2026-05-06T12:10:00Z').getTime(),
clientId: client.id,
shares: 32,
acceptedCount: 1,
address: client.address,
clientName: client.clientName,
sessionId: client.sessionId
});
expect(statistics.hashRate).toBeGreaterThan(0);
});
it('should not suggest a difficulty change before enough time or shares have passed', () => {
+3 -68
View File
@@ -1,4 +1,3 @@
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
import { ClientEntity } from '../ORM/client/client.entity';
const CACHE_SIZE = 30;
@@ -8,31 +7,16 @@ export class StratumV1ClientStatistics {
public targetSubmitShareEveryNSeconds: number = 30;
public hashRate = 0;
private shares: number = 0;
private acceptedCount: number = 0;
private submissionCacheStart: Date;
private submissionCache: { time: Date, difficulty: number }[] = [];
private submissionCacheDifficultySum = 0;
private currentTimeSlot: number = null;
constructor(
private readonly clientStatisticsService: ClientStatisticsService,
) {
constructor() {
this.submissionCacheStart = new Date();
}
// We don't want to save them here because it can be DB intensive, instead do it every once in
// awhile with saveShares()
public async addShares(client: ClientEntity, targetDifficulty: number) {
// 10 min
var coeff = 1000 * 60 * 10;
public async addShares(_client: ClientEntity, targetDifficulty: number) {
var date = new Date();
var timeSlot = new Date(Math.floor(date.getTime() / coeff) * coeff).getTime();
if (this.submissionCache.length > CACHE_SIZE) {
this.submissionCacheDifficultySum -= this.submissionCache[0].difficulty;
@@ -44,55 +28,6 @@ export class StratumV1ClientStatistics {
});
this.submissionCacheDifficultySum += 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);
@@ -151,4 +86,4 @@ export class StratumV1ClientStatistics {
return res;
}
}
}
+63 -13
View File
@@ -107,7 +107,51 @@ describe('StratumV2Client extended channels', () => {
expect(error.errorCode).toBe('invalid-extranonce-size');
});
async function createClient(): Promise<{ client: StratumV2Client; sentFrames: any[] }> {
it('records accepted SV2 shares before presence updates', async () => {
const { client, shareAccountingService, redisMessagingService, jobTemplate } = await createClient();
(client as any).address = 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
(client as any).workerName = 'worker';
(client as any).sessionId = 'sv2-session';
(client as any).userAgent = 'test/sv2';
await (client as any).recordAcceptedShare(
2048,
1024,
jobTemplate,
null,
{
jobId: '1',
nonce: 123,
ntime: parseInt(MockRecording1.TIME, 16),
version: jobTemplate.block.version,
extraNonce2: 'c708000000000000',
},
);
expect(shareAccountingService.recordAcceptedShare).toHaveBeenCalledWith(expect.objectContaining({
protocol: 'sv2',
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
clientName: 'worker',
sessionId: 'sv2-session',
jobId: '1',
creditedDifficulty: 1024,
submissionDifficulty: 2048,
nonce: 123,
extraNonce2: 'c708000000000000',
}));
const accountingCallOrder = shareAccountingService.recordAcceptedShare.mock.invocationCallOrder[0];
const postAccountingPresenceUpdates = redisMessagingService.setClientPresence.mock.invocationCallOrder
.filter(callOrder => callOrder > accountingCallOrder);
expect(postAccountingPresenceUpdates.length).toBeGreaterThan(0);
});
async function createClient(): Promise<{
client: StratumV2Client;
sentFrames: any[];
shareAccountingService: { recordAcceptedShare: jest.Mock };
redisMessagingService: { setClientPresence: jest.Mock; removeClientPresence: jest.Mock };
jobTemplate: any;
}> {
const blockTemplate$ = new BehaviorSubject(MockRecording1.BLOCK_TEMPLATE);
const bitcoinRpcService = {
newBlockTemplate$: blockTemplate$.asObservable(),
@@ -115,7 +159,7 @@ describe('StratumV2Client extended channels', () => {
SUBMIT_BLOCK: jest.fn().mockResolvedValue(null),
};
const stratumV1JobsService = new StratumV1JobsService(bitcoinRpcService as any);
await firstValueFrom(stratumV1JobsService.newMiningJob$);
const jobTemplate = await firstValueFrom(stratumV1JobsService.newMiningJob$);
const mockSocket: any = {
destroyed: false,
@@ -138,11 +182,24 @@ describe('StratumV2Client extended channels', () => {
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
clientName: 'worker',
userAgent: 'test/sv2',
startTime: new Date(),
bestDifficulty: 0,
} as unknown as ClientEntity;
let nextChannelId = 1;
const sentFrames: any[] = [];
const clientService = {
insert: jest.fn().mockResolvedValue(clientEntity),
delete: jest.fn().mockResolvedValue(undefined),
updateBestDifficultyIfHigher: jest.fn().mockResolvedValue(undefined),
};
const shareAccountingService = {
recordAcceptedShare: jest.fn().mockResolvedValue(undefined),
};
const redisMessagingService = {
setClientPresence: jest.fn().mockResolvedValue(undefined),
removeClientPresence: jest.fn().mockResolvedValue(undefined),
};
const client = new StratumV2Client(
socket,
Buffer.alloc(0),
@@ -167,16 +224,7 @@ describe('StratumV2Client extended channels', () => {
} as any,
stratumV1JobsService,
bitcoinRpcService as any,
{
insert: jest.fn().mockResolvedValue(clientEntity),
delete: jest.fn().mockResolvedValue(undefined),
heartbeatBulkAsync: jest.fn(),
updateBestDifficultyIfHigher: jest.fn().mockResolvedValue(undefined),
} as any,
{
insert: jest.fn().mockResolvedValue(undefined),
updateBulkAsync: jest.fn(),
} as any,
clientService as any,
{ notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined) } as any,
{ save: jest.fn().mockResolvedValue(undefined) } as any,
{
@@ -195,12 +243,14 @@ describe('StratumV2Client extended channels', () => {
resetBestDifficultyAndShares: jest.fn().mockResolvedValue(undefined),
updateBestDifficultyIfHigher: jest.fn().mockResolvedValue(undefined),
} as any,
shareAccountingService as any,
redisMessagingService as any,
);
(client as any).sendFrame = jest.fn((msgType: number, payload: Buffer, extensionType = 0) => {
sentFrames.push({ msgType, payload, extensionType });
return Promise.resolve();
});
return { client, sentFrames };
return { client, sentFrames, shareAccountingService, redisMessagingService, jobTemplate };
}
});
+70 -9
View File
@@ -7,11 +7,12 @@ import { firstValueFrom, Subscription } from 'rxjs';
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
import { BlocksService } from '../ORM/blocks/blocks.service';
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
import { ClientEntity } from '../ORM/client/client.entity';
import { ClientService } from '../ORM/client/client.service';
import { ShareAccountingService } from '../ORM/share-accounting/share-accounting.service';
import { BitcoinRpcService } from '../services/bitcoin-rpc.service';
import { NotificationService } from '../services/notification.service';
import { RedisMessagingService } from '../services/redis-messaging.service';
import { StratumV2Service } from '../services/stratum-v2.service';
import { IJobTemplate, StratumV1JobsService } from '../services/stratum-v1-jobs.service';
import { patchCoinbasePrefixVarint } from '../utils/coinbase-prefix.utils';
@@ -128,18 +129,19 @@ export class StratumV2Client {
private readonly stratumV1JobsService: StratumV1JobsService,
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly clientService: ClientService,
private readonly clientStatisticsService: ClientStatisticsService,
private readonly notificationService: NotificationService,
private readonly blocksService: BlocksService,
private readonly configService: ConfigService,
private readonly addressSettingsService: AddressSettingsService,
private readonly shareAccountingService?: ShareAccountingService,
private readonly redisMessagingService?: RedisMessagingService,
) {
this.firstChunkSummary = this.describeChunk(firstChunk);
this.noiseSession = new Sv2NoiseSession(this.stratumV2Service.getNoiseConfig());
this.sessionDifficulty = this.getInitialDifficulty();
this.targetSharesPerMinute = this.getTargetSharesPerMinute();
this.difficultyCheckIntervalMs = this.getDifficultyCheckIntervalMs();
this.statistics = new StratumV1ClientStatistics(this.clientStatisticsService);
this.statistics = new StratumV1ClientStatistics();
this.statistics.targetSubmitShareEveryNSeconds = 60 / this.targetSharesPerMinute;
this.network = this.getNetwork();
@@ -171,6 +173,7 @@ export class StratumV2Client {
}
this.channels.clear();
if (this.clientEntity?.id != null) {
await this.redisMessagingService?.removeClientPresence(this.clientEntity.id, this.clientEntity.address);
await this.clientService.delete(this.clientEntity.id);
}
}
@@ -647,7 +650,13 @@ export class StratumV2Client {
if (submissionDifficulty >= extendedJob.jobTemplate.blockData.networkDifficulty) {
updatedJobBlock = this.reconstructExtendedBlock(extendedJob, submission, merkleRoot, channel.extranoncePrefix);
}
await this.recordAcceptedShare(submissionDifficulty, jobDifficulty, extendedJob.jobTemplate, updatedJobBlock);
await this.recordAcceptedShare(submissionDifficulty, jobDifficulty, extendedJob.jobTemplate, updatedJobBlock, {
jobId: submission.jobId.toString(16),
nonce: submission.nonce,
ntime: submission.ntime,
version: submission.version,
extraNonce2: submission.extranonce.toString('hex'),
});
}
private async handleAcceptedShare(
@@ -671,7 +680,13 @@ export class StratumV2Client {
);
}
await this.recordAcceptedShare(submissionDifficulty, jobDifficulty, jobTemplate, updatedJobBlock);
await this.recordAcceptedShare(submissionDifficulty, jobDifficulty, jobTemplate, updatedJobBlock, {
jobId: job.jobId,
nonce: submission.nonce,
ntime: submission.ntime,
version: submission.version,
extraNonce2: FIXED_STANDARD_EXTRANONCE2,
});
}
private async recordAcceptedShare(
@@ -679,13 +694,15 @@ export class StratumV2Client {
jobDifficulty: number,
jobTemplate: IJobTemplate,
updatedJobBlock: bitcoinjs.Block | null,
share: { jobId: string; nonce: number; ntime: number; version: number; extraNonce2: string },
): Promise<void> {
await this.ensureClientEntity();
let blockSubmissionResult: string = null;
if (updatedJobBlock != null) {
console.log(`[SV2 ${this.sessionId}] BLOCK FOUND at height ${jobTemplate.blockData.height}`);
const blockHex = updatedJobBlock.toHex(false);
const result = await this.bitcoinRpcService.SUBMIT_BLOCK(blockHex);
blockSubmissionResult = await this.bitcoinRpcService.SUBMIT_BLOCK(blockHex);
await this.blocksService.save({
height: jobTemplate.blockData.height,
minerAddress: this.address,
@@ -697,17 +714,37 @@ export class StratumV2Client {
this.address,
jobTemplate.blockData.height,
updatedJobBlock,
result,
blockSubmissionResult,
);
if (result == null) {
if (blockSubmissionResult == null) {
await this.addressSettingsService.resetBestDifficultyAndShares();
}
}
await this.shareAccountingService?.recordAcceptedShare({
protocol: 'sv2',
address: this.address,
clientName: this.workerName,
sessionId: this.sessionId,
clientId: this.clientEntity.id,
jobId: share.jobId,
jobTemplateId: jobTemplate.blockData.id,
blockHeight: jobTemplate.blockData.height,
creditedDifficulty: jobDifficulty,
submissionDifficulty,
networkDifficulty: jobTemplate.blockData.networkDifficulty,
nonce: share.nonce,
ntime: share.ntime,
version: share.version,
extraNonce2: share.extraNonce2,
isBlockCandidate: updatedJobBlock != null,
blockSubmissionResult,
});
await this.statistics.addShares(this.clientEntity, jobDifficulty);
const now = new Date();
this.clientService.heartbeatBulkAsync(this.clientEntity.id, this.statistics.hashRate, now);
this.clientEntity.updatedAt = now;
this.clientEntity.hashRate = this.statistics.hashRate;
await this.updateClientPresence(now);
if (submissionDifficulty > this.clientEntity.bestDifficulty) {
await this.clientService.updateBestDifficultyIfHigher(this.clientEntity.id, submissionDifficulty);
@@ -717,6 +754,7 @@ export class StratumV2Client {
submissionDifficulty,
this.userAgent,
);
await this.updateClientPresence(new Date());
}
}
@@ -1085,12 +1123,35 @@ export class StratumV2Client {
startTime: new Date(),
bestDifficulty: 0,
});
await this.updateClientPresence(new Date());
})();
}
await this.creatingEntity;
}
private async updateClientPresence(lastSeen: Date): Promise<void> {
if (this.clientEntity == null) {
return;
}
try {
await this.redisMessagingService?.setClientPresence({
clientId: this.clientEntity.id,
address: this.clientEntity.address,
clientName: this.clientEntity.clientName,
sessionId: this.clientEntity.sessionId,
userAgent: this.clientEntity.userAgent,
startTime: new Date(this.clientEntity.startTime).toISOString(),
lastSeen: lastSeen.toISOString(),
hashRate: this.statistics.hashRate,
bestDifficulty: Number(this.clientEntity.bestDifficulty ?? 0),
});
} catch (error) {
console.error(`Failed to update SV2 client presence: ${error.message}`);
}
}
private parseUserIdentity(userIdentity: string): { address: string; workerName: string } {
const parts = userIdentity.split('.');
const rawAddress = parts[0] ?? '';
+16
View File
@@ -0,0 +1,16 @@
import { AppDataSource } from '../data-source';
async function main() {
await AppDataSource.initialize();
try {
const migrations = await AppDataSource.runMigrations();
console.log(`Ran ${migrations.length} database migrations`);
} finally {
await AppDataSource.destroy();
}
}
main().catch(error => {
console.error('Database migration failed:', error);
process.exit(1);
});
+68
View File
@@ -0,0 +1,68 @@
import { spawn } from 'child_process';
import * as net from 'net';
async function main() {
await waitForTcp(process.env.DB_HOST, parseInt(process.env.DB_PORT ?? '5432', 10), 'TimescaleDB');
const redisUrl = new URL(process.env.REDIS_URL ?? 'redis://redis:6379');
await waitForTcp(redisUrl.hostname, parseInt(redisUrl.port || '6379', 10), 'Redis');
await run('node', ['dist/scripts/run-migrations.js']);
if (process.env.PM2_ENABLED === 'false') {
await run('node', ['dist/main.js'], true);
return;
}
await run('./node_modules/.bin/pm2-runtime', ['ecosystem.config.js'], true);
}
function waitForTcp(host: string, port: number, label: string): Promise<void> {
const deadline = Date.now() + parseInt(process.env.STARTUP_WAIT_TIMEOUT_MS ?? '60000', 10);
return new Promise((resolve, reject) => {
const attempt = () => {
const socket = net.createConnection({ host, port });
socket.once('connect', () => {
socket.destroy();
console.log(`${label} is reachable at ${host}:${port}`);
resolve();
});
socket.once('error', () => {
socket.destroy();
if (Date.now() > deadline) {
reject(new Error(`${label} was not reachable at ${host}:${port}`));
return;
}
setTimeout(attempt, 1000);
});
};
attempt();
});
}
function run(command: string, args: string[], inherit = false): Promise<void> {
return new Promise((resolve, reject) => {
const child = spawn(command, args, {
stdio: inherit ? 'inherit' : 'pipe',
env: process.env,
});
if (!inherit) {
child.stdout?.on('data', data => process.stdout.write(data));
child.stderr?.on('data', data => process.stderr.write(data));
}
child.once('exit', code => {
if (code === 0) {
resolve();
return;
}
reject(new Error(`${command} ${args.join(' ')} exited with ${code}`));
});
});
}
main().catch(error => {
console.error(error);
process.exit(1);
});
+4 -70
View File
@@ -1,21 +1,15 @@
import { Injectable, OnModuleInit } from '@nestjs/common';
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 { 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()
export class AppService implements OnModuleInit {
constructor(
private readonly clientStatisticsService: ClientStatisticsService,
private readonly clientService: ClientService,
private readonly rpcBlockService: RpcBlockService,
private readonly homeGraphService: HomeGraphService,
private readonly dataSource: DataSource,
private readonly userAgentReportService: UserAgentReportService
) {
@@ -25,26 +19,14 @@ export class AppService implements OnModuleInit {
if (process.env.MASTER == 'true') {
setInterval(async () => {
await this.deleteOldStatistics();
await this.deleteOldClients();
}, 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();
@@ -52,61 +34,13 @@ export class AppService implements OnModuleInit {
}, 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);
}
private async deleteOldStatistics() {
console.log('Deleting statistics');
private async deleteOldClients() {
console.log('Deleting old clients');
const deletedStatistics = await this.clientStatisticsService.deleteOldStatistics();
console.log(`Deleted ${deletedStatistics.affected} old statistics`);
const deletedClients = await this.clientService.deleteOldClients();
console.log(`Deleted ${deletedClients.affected} old clients`);
}
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);
}
}
}
+36 -11
View File
@@ -7,7 +7,7 @@ import * as zmq from 'zeromq';
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo';
import * as PGPubsub from 'pg-pubsub';
import { RedisMessagingService } from './redis-messaging.service';
@Injectable()
export class BitcoinRpcService implements OnModuleInit {
@@ -15,7 +15,6 @@ export class BitcoinRpcService implements OnModuleInit {
private client: AxiosInstance;
private _newBlockTemplate$: BehaviorSubject<IBlockTemplate> = new BehaviorSubject(undefined);
private pubsubInstance: PGPubsub;
private resetTemplateInterval$ = new Subject<void>();
private rpcRequestId = 0;
@@ -24,16 +23,14 @@ export class BitcoinRpcService implements OnModuleInit {
constructor(
private readonly configService: ConfigService,
private rpcBlockService: RpcBlockService
private rpcBlockService: RpcBlockService,
private readonly redisMessagingService: RedisMessagingService
) {
}
async onModuleInit() {
this.pubsubInstance = new PGPubsub('postgres://' + this.configService.get('DB_USERNAME') + ':' + this.configService.get('DB_PASSWORD') + '@' + this.configService.get('DB_HOST') + ':' + this.configService.get('DB_PORT') + '/' + this.configService.get('DB_DATABASE'))
const url = this.configService.get('BITCOIN_RPC_URL');
const user = this.configService.get('BITCOIN_RPC_USER');
const pass = this.configService.get('BITCOIN_RPC_PASSWORD');
@@ -60,11 +57,10 @@ export class BitcoinRpcService implements OnModuleInit {
console.log(`MASTER? ${process.env.MASTER}`)
if (process.env.MASTER != 'true') {
this.pubsubInstance.addChannel('miningInfo', async (miningInfo: IMiningInfo) => {
//console.log('PG Sub. new template');
await this.loadLatestTemplateForWorker();
await this.redisMessagingService.subscribeMiningInfoUpdates(async (miningInfo: IMiningInfo) => {
this.miningInfo = miningInfo;
const savedBlockTemplate = await this.rpcBlockService.getSavedBlockTemplate(miningInfo.blocks);
this._newBlockTemplate$.next(JSON.parse(savedBlockTemplate.data));
await this.loadTemplateForWorker(miningInfo.blocks);
});
} else {
console.log('Using ZMQ');
@@ -109,7 +105,36 @@ export class BitcoinRpcService implements OnModuleInit {
public async getAndBroadcastLatestTemplate() {
const blockTemplate = await this.loadBlockTemplate(this.miningInfo.blocks);
this._newBlockTemplate$.next(blockTemplate);
await this.pubsubInstance.publish('miningInfo', this.miningInfo);
await this.redisMessagingService.setLatestMiningInfo(this.miningInfo);
await this.redisMessagingService.setBlockTemplate(this.miningInfo.blocks, blockTemplate);
await this.redisMessagingService.publishMiningInfoUpdate(this.miningInfo);
}
private async loadLatestTemplateForWorker() {
const latestMiningInfo = await this.redisMessagingService.getLatestMiningInfo();
if (latestMiningInfo != null) {
this.miningInfo = latestMiningInfo;
await this.loadTemplateForWorker(latestMiningInfo.blocks);
return;
}
const latestBlockTemplate = await this.redisMessagingService.getLatestBlockTemplate();
if (latestBlockTemplate != null) {
this._newBlockTemplate$.next(latestBlockTemplate);
}
}
private async loadTemplateForWorker(blockHeight: number) {
const redisBlockTemplate = await this.redisMessagingService.getBlockTemplate(blockHeight);
if (redisBlockTemplate != null) {
this._newBlockTemplate$.next(redisBlockTemplate);
return;
}
const savedBlockTemplate = await this.rpcBlockService.getSavedBlockTemplate(blockHeight);
if (savedBlockTemplate?.data != null) {
this._newBlockTemplate$.next(JSON.parse(savedBlockTemplate.data));
}
}
private async loadBlockTemplate(blockHeight: number) {
+12
View File
@@ -0,0 +1,12 @@
import { Global, Module } from '@nestjs/common';
import { ConfigModule } from '@nestjs/config';
import { RedisMessagingService } from './redis-messaging.service';
@Global()
@Module({
imports: [ConfigModule],
providers: [RedisMessagingService],
exports: [RedisMessagingService],
})
export class RedisMessagingModule { }
@@ -0,0 +1,182 @@
import { ConfigService } from '@nestjs/config';
import { createClient } from 'redis';
import { RedisMessagingService } from './redis-messaging.service';
jest.mock('redis', () => ({
createClient: jest.fn(),
}));
describe('RedisMessagingService', () => {
let clients: any[];
let service: RedisMessagingService;
beforeEach(() => {
store.clear();
sets.clear();
clientsByRole.publisher = null;
clientsByRole.subscriber = null;
clients = [createRedisClient(), createRedisClient()];
(createClient as jest.Mock).mockImplementation(() => clients.shift());
service = new RedisMessagingService({
get: jest.fn((key: string) => key === 'REDIS_URL' ? 'redis://test-redis:6379' : null),
} as unknown as ConfigService);
});
afterEach(() => {
jest.clearAllMocks();
});
it('should store and replay latest mining info and block templates', async () => {
await service.connect();
await service.setLatestMiningInfo({ blocks: 900000 } as any);
await service.setBlockTemplate(900000, { height: 900000, transactions: [] } as any);
expect(await service.getLatestMiningInfo()).toEqual({ blocks: 900000 });
expect(await service.getBlockTemplate(900000)).toEqual({ height: 900000, transactions: [] });
expect(await service.getLatestBlockTemplate()).toEqual({ height: 900000, transactions: [] });
});
it('should publish and subscribe to mining info updates', async () => {
await service.connect();
const handler = jest.fn().mockResolvedValue(undefined);
await service.subscribeMiningInfoUpdates(handler);
await service.publishMiningInfoUpdate({ blocks: 900001 } as any);
expect(handler).toHaveBeenCalledWith({ blocks: 900001 });
});
it('should ignore malformed pubsub messages', async () => {
await service.connect();
const handler = jest.fn();
const consoleSpy = jest.spyOn(console, 'error').mockImplementation(() => undefined);
await service.subscribeMiningInfoUpdates(handler);
await clientsByRole.subscriber.callback('{bad json');
expect(handler).not.toHaveBeenCalled();
expect(consoleSpy).toHaveBeenCalledWith(expect.stringContaining('Invalid Redis mining info update'));
consoleSpy.mockRestore();
});
it('should store, index, and remove client presence', async () => {
await service.connect();
await service.setClientPresence({
clientId: '00000000-0000-4000-8000-000000000001',
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
clientName: 'worker',
sessionId: '57a6f098',
userAgent: 'bitaxe',
startTime: '2026-06-07T12:00:00.000Z',
lastSeen: '2026-06-07T12:01:00.000Z',
hashRate: 100,
bestDifficulty: 200,
});
expect(await service.getClientPresence('00000000-0000-4000-8000-000000000001'))
.toEqual(expect.objectContaining({
clientId: '00000000-0000-4000-8000-000000000001',
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
hashRate: 100,
bestDifficulty: 200,
}));
expect(await service.getClientPresenceByAddress('tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4'))
.toHaveLength(1);
expect(await service.getAllClientPresence()).toHaveLength(1);
await service.removeClientPresence('00000000-0000-4000-8000-000000000001');
expect(await service.getAllClientPresence()).toHaveLength(0);
});
it('should clear all client presence keys', async () => {
await service.connect();
await service.setClientPresence({
clientId: '00000000-0000-4000-8000-000000000001',
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
clientName: 'worker',
sessionId: '57a6f098',
userAgent: 'bitaxe',
startTime: '2026-06-07T12:00:00.000Z',
lastSeen: '2026-06-07T12:01:00.000Z',
hashRate: 100,
bestDifficulty: 200,
});
await service.clearClientPresence();
expect(await service.getAllClientPresence()).toHaveLength(0);
expect([...store.keys()].filter(key => key.startsWith('client-presence'))).toHaveLength(0);
expect([...sets.keys()].filter(key => key.startsWith('client-presence'))).toHaveLength(0);
});
});
const store = new Map<string, string>();
const sets = new Map<string, Set<string>>();
const clientsByRole: { publisher?: any; subscriber?: any } = {};
function createRedisClient() {
const client = {
connect: jest.fn().mockResolvedValue(undefined),
quit: jest.fn().mockResolvedValue(undefined),
on: jest.fn(),
set: jest.fn((key: string, value: string) => {
store.set(key, value);
return Promise.resolve('OK');
}),
setEx: jest.fn((key: string, seconds: number, value: string) => {
store.set(key, value);
return Promise.resolve('OK');
}),
get: jest.fn((key: string) => Promise.resolve(store.get(key) ?? null)),
del: jest.fn((...args: (string | string[])[]) => {
const keys = args.flatMap(key => Array.isArray(key) ? key : [key]);
let deleted = 0;
keys.forEach(item => {
deleted += store.delete(item) ? 1 : 0;
deleted += sets.delete(item) ? 1 : 0;
});
return Promise.resolve(deleted);
}),
sAdd: jest.fn((key: string, value: string) => {
const set = sets.get(key) ?? new Set<string>();
set.add(value);
sets.set(key, set);
return Promise.resolve(1);
}),
sRem: jest.fn((key: string, value: string) => {
const set = sets.get(key);
const deleted = set?.delete(value) ? 1 : 0;
return Promise.resolve(deleted);
}),
sMembers: jest.fn((key: string) => Promise.resolve([...sets.get(key) ?? []])),
scanIterator: jest.fn(async function* scanIterator({ MATCH }: { MATCH: string }) {
const prefix = MATCH.replace('*', '');
for (const key of [...store.keys(), ...sets.keys()]) {
if (key.startsWith(prefix)) {
yield key;
}
}
}),
publish: jest.fn((channel: string, value: string) => {
void clientsByRole.subscriber?.callback(value);
return Promise.resolve(1);
}),
subscribe: jest.fn((channel: string, callback: (message: string) => Promise<void>) => {
clientsByRole.subscriber.callback = callback;
return Promise.resolve();
}),
};
if (clientsByRole.publisher == null) {
clientsByRole.publisher = client;
} else {
clientsByRole.subscriber = client;
}
return client;
}
+291
View File
@@ -0,0 +1,291 @@
import { Injectable, OnModuleDestroy, OnModuleInit } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { createClient, RedisClientType } from 'redis';
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo';
const MINING_INFO_CHANNEL = 'mining-info.updated';
const MINING_INFO_KEY = 'mining-info:latest';
const BLOCK_TEMPLATE_LATEST_KEY = 'block-template:latest';
const CLIENT_PRESENCE_TTL_SECONDS = 180;
const blockTemplateKey = (height: number) => `block-template:${height}`;
const CLIENT_PRESENCE_ALL_KEY = 'client-presence:all';
const clientPresenceKey = (clientId: string) => `client-presence:${clientId}`;
const clientPresenceAddressKey = (address: string) => `client-presence:address:${address}`;
export interface ClientPresence {
clientId: string;
address: string;
clientName: string;
sessionId: string;
userAgent?: string | null;
startTime: string;
lastSeen: string;
hashRate: number;
bestDifficulty: number;
}
@Injectable()
export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
private publisher: RedisClientType;
private subscriber: RedisClientType;
private connected = false;
private readonly clientPresenceTtlSeconds: number;
constructor(
private readonly configService: ConfigService,
) {
this.clientPresenceTtlSeconds = this.readPositiveInt('CLIENT_PRESENCE_TTL_SECONDS', CLIENT_PRESENCE_TTL_SECONDS);
}
public async onModuleInit() {
await this.connect().catch(error => {
console.error(`Redis unavailable at startup: ${error.message}`);
});
}
public async onModuleDestroy() {
await Promise.all([
this.publisher?.quit().catch(() => undefined),
this.subscriber?.quit().catch(() => undefined),
]);
}
public async connect() {
if (this.connected) {
return;
}
const url = this.configService.get<string>('REDIS_URL') ?? process.env.REDIS_URL ?? 'redis://localhost:6379';
const socket = {
reconnectStrategy: (retries: number) => Math.min(retries * 250, 5000),
};
this.publisher = createClient({ url, socket });
this.subscriber = createClient({ url, socket });
this.publisher.on('error', error => console.error(`Redis publisher error: ${error.message}`));
this.subscriber.on('error', error => console.error(`Redis subscriber error: ${error.message}`));
this.publisher.on('end', () => { this.connected = false; });
this.subscriber.on('end', () => { this.connected = false; });
await Promise.all([this.publisher.connect(), this.subscriber.connect()]);
this.connected = true;
}
public async publishMiningInfoUpdate(miningInfo: IMiningInfo) {
if (!await this.ensureConnected()) {
return;
}
await this.publisher.publish(MINING_INFO_CHANNEL, JSON.stringify(miningInfo));
}
public async subscribeMiningInfoUpdates(handler: (miningInfo: IMiningInfo) => Promise<void>) {
if (!await this.ensureConnected()) {
return;
}
await this.subscriber.subscribe(MINING_INFO_CHANNEL, async message => {
try {
await handler(JSON.parse(message));
} catch (error) {
console.error(`Invalid Redis mining info update: ${error.message}`);
}
});
}
public async setLatestMiningInfo(miningInfo: IMiningInfo) {
if (!await this.ensureConnected()) {
return;
}
await this.publisher.set(MINING_INFO_KEY, JSON.stringify(miningInfo));
}
public async getLatestMiningInfo(): Promise<IMiningInfo | null> {
if (!await this.ensureConnected()) {
return null;
}
const value = await this.publisher.get(MINING_INFO_KEY);
return value == null ? null : JSON.parse(value as string);
}
public async setBlockTemplate(height: number, blockTemplate: IBlockTemplate) {
if (!await this.ensureConnected()) {
return;
}
const serialized = JSON.stringify(blockTemplate);
await Promise.all([
this.publisher.set(blockTemplateKey(height), serialized),
this.publisher.set(BLOCK_TEMPLATE_LATEST_KEY, serialized),
]);
}
public async getBlockTemplate(height: number): Promise<IBlockTemplate | null> {
if (!await this.ensureConnected()) {
return null;
}
const value = await this.publisher.get(blockTemplateKey(height));
return value == null ? null : JSON.parse(value as string);
}
public async getLatestBlockTemplate(): Promise<IBlockTemplate | null> {
if (!await this.ensureConnected()) {
return null;
}
const value = await this.publisher.get(BLOCK_TEMPLATE_LATEST_KEY);
return value == null ? null : JSON.parse(value as string);
}
public async setClientPresence(presence: ClientPresence): Promise<void> {
if (!await this.ensureConnected()) {
return;
}
const serialized = JSON.stringify({
...presence,
userAgent: presence.userAgent ?? null,
hashRate: Number.isFinite(Number(presence.hashRate)) ? Number(presence.hashRate) : 0,
bestDifficulty: Number.isFinite(Number(presence.bestDifficulty)) ? Number(presence.bestDifficulty) : 0,
});
await Promise.all([
this.publisher.setEx(clientPresenceKey(presence.clientId), this.clientPresenceTtlSeconds, serialized),
this.publisher.sAdd(CLIENT_PRESENCE_ALL_KEY, presence.clientId),
this.publisher.sAdd(clientPresenceAddressKey(presence.address), presence.clientId),
]);
}
public async removeClientPresence(clientId: string, address?: string): Promise<void> {
if (!await this.ensureConnected()) {
return;
}
let resolvedAddress = address;
if (resolvedAddress == null) {
const presence = await this.getClientPresence(clientId);
resolvedAddress = presence?.address;
}
const removals: Promise<unknown>[] = [
this.publisher.del(clientPresenceKey(clientId)),
this.publisher.sRem(CLIENT_PRESENCE_ALL_KEY, clientId),
];
if (resolvedAddress != null) {
removals.push(this.publisher.sRem(clientPresenceAddressKey(resolvedAddress), clientId));
}
await Promise.all(removals);
}
public async getClientPresence(clientId: string): Promise<ClientPresence | null> {
if (!await this.ensureConnected()) {
return null;
}
const value = await this.publisher.get(clientPresenceKey(clientId));
return this.parseClientPresence(value);
}
public async getClientPresenceByAddress(address: string): Promise<ClientPresence[]> {
if (!await this.ensureConnected()) {
return [];
}
return this.getPresenceFromSet(clientPresenceAddressKey(address));
}
public async getAllClientPresence(): Promise<ClientPresence[]> {
if (!await this.ensureConnected()) {
return [];
}
return this.getPresenceFromSet(CLIENT_PRESENCE_ALL_KEY);
}
public async clearClientPresence(): Promise<void> {
if (!await this.ensureConnected()) {
return;
}
const batch: string[] = [];
for await (const keyOrKeys of (this.publisher as any).scanIterator({ MATCH: 'client-presence*', COUNT: 1000 })) {
const keys = Array.isArray(keyOrKeys) ? keyOrKeys : [keyOrKeys];
batch.push(...keys.map(key => key as string));
if (batch.length >= 500) {
await this.deleteKeys(batch.splice(0));
}
}
await this.deleteKeys(batch);
}
private async ensureConnected(): Promise<boolean> {
if (!this.connected) {
try {
await this.connect();
} catch (error) {
console.error(`Redis messaging degraded: ${error.message}`);
return false;
}
}
return true;
}
private async getPresenceFromSet(setKey: string): Promise<ClientPresence[]> {
const clientIds = await this.publisher.sMembers(setKey);
if (clientIds.length === 0) {
return [];
}
const presences = await Promise.all(clientIds.map(async clientId => {
const presence = await this.getClientPresence(clientId);
if (presence == null) {
await this.publisher.sRem(setKey, clientId);
if (setKey !== CLIENT_PRESENCE_ALL_KEY) {
await this.publisher.sRem(CLIENT_PRESENCE_ALL_KEY, clientId);
}
}
return presence;
}));
return presences.filter((presence): presence is ClientPresence => presence != null);
}
private parseClientPresence(value: unknown): ClientPresence | null {
if (value == null) {
return null;
}
try {
const parsed = JSON.parse(value as string);
if (parsed?.clientId == null || parsed?.address == null) {
return null;
}
return {
clientId: parsed.clientId,
address: parsed.address,
clientName: parsed.clientName ?? 'default',
sessionId: parsed.sessionId ?? parsed.clientId,
userAgent: parsed.userAgent ?? null,
startTime: parsed.startTime,
lastSeen: parsed.lastSeen,
hashRate: Number(parsed.hashRate ?? 0),
bestDifficulty: Number(parsed.bestDifficulty ?? 0),
};
} catch (error) {
console.error(`Invalid Redis client presence: ${error.message}`);
return null;
}
}
private readPositiveInt(name: string, defaultValue: number): number {
const value = Number(this.configService.get<string>(name) ?? process.env[name]);
return Number.isInteger(value) && value > 0 ? value : defaultValue;
}
private async deleteKeys(keys: string[]): Promise<void> {
if (keys.length === 0) {
return;
}
await (this.publisher as any).del(...keys);
}
}
+15 -2
View File
@@ -11,7 +11,9 @@ describe('StratumV1Service', () => {
let service: StratumV1Service;
let clientService;
let userAgentReportService;
let stratumV2Service;
let redisMessagingService;
let consoleLogSpy: jest.SpyInstance;
let consoleWarnSpy: jest.SpyInstance;
@@ -20,10 +22,16 @@ describe('StratumV1Service', () => {
clientService = {
deleteAll: jest.fn().mockResolvedValue(undefined)
};
userAgentReportService = {
refreshReport: jest.fn().mockResolvedValue(undefined)
};
stratumV2Service = {
ensureInitialized: jest.fn().mockResolvedValue(undefined),
createClient: jest.fn()
};
redisMessagingService = {
clearClientPresence: jest.fn().mockResolvedValue(undefined)
};
service = new StratumV1Service(
{} as any,
clientService,
@@ -32,8 +40,10 @@ describe('StratumV1Service', () => {
{} as any,
{} as any,
{} as any,
{} as any,
stratumV2Service as any
stratumV2Service as any,
userAgentReportService as any,
undefined,
redisMessagingService as any
);
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
consoleWarnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined);
@@ -61,6 +71,8 @@ describe('StratumV1Service', () => {
jest.runOnlyPendingTimers();
expect(clientService.deleteAll).toHaveBeenCalled();
expect(redisMessagingService.clearClientPresence).toHaveBeenCalled();
expect(userAgentReportService.refreshReport).toHaveBeenCalled();
expect(startSocketServerSpy).not.toHaveBeenCalled();
expect(startSecureSocketServerSpy).not.toHaveBeenCalled();
expect(consoleLogSpy).toHaveBeenCalledWith('Master process skipping Stratum socket listeners');
@@ -78,6 +90,7 @@ describe('StratumV1Service', () => {
jest.advanceTimersByTime(10000);
expect(clientService.deleteAll).not.toHaveBeenCalled();
expect(userAgentReportService.refreshReport).not.toHaveBeenCalled();
expect(startSocketServerSpy).toHaveBeenCalledWith(3333);
expect(startSocketServerSpy).toHaveBeenCalledWith(3334);
expect(startSecureSocketServerSpy).toHaveBeenCalledWith(4333);
+12 -5
View File
@@ -5,12 +5,14 @@ import { monitorEventLoopDelay } from 'perf_hooks';
import { StratumV1Client } from '../models/StratumV1Client';
import { StratumV2Client } from '../models/StratumV2Client';
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 { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
import { ClientService } from '../ORM/client/client.service';
import { ShareAccountingService } from '../ORM/share-accounting/share-accounting.service';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { NotificationService } from './notification.service';
import { RedisMessagingService } from './redis-messaging.service';
import { StratumV1JobsService } from './stratum-v1-jobs.service';
import { StratumV2Service } from './stratum-v2.service';
@@ -51,13 +53,15 @@ export class StratumV1Service implements OnModuleInit {
constructor(
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly clientService: ClientService,
private readonly clientStatisticsService: ClientStatisticsService,
private readonly notificationService: NotificationService,
private readonly blocksService: BlocksService,
private readonly configService: ConfigService,
private readonly stratumV1JobsService: StratumV1JobsService,
private readonly addressSettingsService: AddressSettingsService,
private readonly stratumV2Service: StratumV2Service
private readonly stratumV2Service: StratumV2Service,
private readonly userAgentReportService: UserAgentReportService,
private readonly shareAccountingService?: ShareAccountingService,
private readonly redisMessagingService?: RedisMessagingService
) {
}
@@ -66,6 +70,8 @@ export class StratumV1Service implements OnModuleInit {
if (process.env.MASTER == 'true') {
await this.clientService.deleteAll();
await this.redisMessagingService?.clearClientPresence();
await this.userAgentReportService.refreshReport();
console.log('Master process skipping Stratum socket listeners');
return;
}
@@ -191,11 +197,12 @@ export class StratumV1Service implements OnModuleInit {
this.stratumV1JobsService,
this.bitcoinRpcService,
this.clientService,
this.clientStatisticsService,
this.notificationService,
this.blocksService,
this.configService,
this.addressSettingsService
this.addressSettingsService,
this.shareAccountingService,
this.redisMessagingService
);
}
+6 -3
View File
@@ -5,8 +5,8 @@ import { Server, Socket } from 'net';
import { AddressSettingsService } from '../ORM/address-settings/address-settings.service';
import { BlocksService } from '../ORM/blocks/blocks.service';
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
import { ClientService } from '../ORM/client/client.service';
import { ShareAccountingService } from '../ORM/share-accounting/share-accounting.service';
import { StratumV2Client } from '../models/StratumV2Client';
import { encodeSv2AuthorityPublicKey } from '../models/sv2/sv2-authority-key';
import { Sv2ExtranonceManager } from '../models/sv2/sv2-extranonce-manager';
@@ -19,6 +19,7 @@ import {
} from '../models/sv2/sv2-noise';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { NotificationService } from './notification.service';
import { RedisMessagingService } from './redis-messaging.service';
import { StratumV1JobsService } from './stratum-v1-jobs.service';
@Injectable()
@@ -35,12 +36,13 @@ export class StratumV2Service implements OnModuleInit {
constructor(
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly clientService: ClientService,
private readonly clientStatisticsService: ClientStatisticsService,
private readonly notificationService: NotificationService,
private readonly blocksService: BlocksService,
private readonly configService: ConfigService,
private readonly stratumV1JobsService: StratumV1JobsService,
private readonly addressSettingsService: AddressSettingsService,
private readonly shareAccountingService?: ShareAccountingService,
private readonly redisMessagingService?: RedisMessagingService,
) {}
public async onModuleInit(): Promise<void> {
@@ -77,11 +79,12 @@ export class StratumV2Service implements OnModuleInit {
this.stratumV1JobsService,
this.bitcoinRpcService,
this.clientService,
this.clientStatisticsService,
this.notificationService,
this.blocksService,
this.configService,
this.addressSettingsService,
this.shareAccountingService,
this.redisMessagingService,
);
}