mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 09:05:06 -07:00
Compare commits
14
Commits
d05d5f0b37
...
deaf3e56bd
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
deaf3e56bd | ||
|
|
9b73e33b13 | ||
|
|
0560b874dc | ||
|
|
7f6f42c5c9 | ||
|
|
4c59797886 | ||
|
|
4c91e69ca8 | ||
|
|
7b7dc87242 | ||
|
|
df35549f27 | ||
|
|
8c795a3fa2 | ||
|
|
ac7827a944 | ||
|
|
17ba4afc1a | ||
|
|
b612b50f3e | ||
|
|
b6a951bdee | ||
|
|
817f416b79 |
@@ -19,6 +19,7 @@ API_PORT=3334
|
||||
# Keep API_BIND_HOST as 127.0.0.1 when a reverse proxy terminates public traffic.
|
||||
API_BIND_HOST=127.0.0.1
|
||||
API_PUBLIC_PORT=3334
|
||||
API_WORKERS=4
|
||||
|
||||
# Plain TCP Stratum ports accept both SV1 JSON-RPC and SV2 Noise/binary traffic.
|
||||
STRATUM_PORTS=3333,3332,3331,3330
|
||||
|
||||
@@ -50,6 +50,7 @@ services:
|
||||
REDIS_URL: ${REDIS_URL:-redis://redis:6379}
|
||||
API_PORT: ${API_PORT:-3334}
|
||||
API_SECURE: ${API_SECURE:-false}
|
||||
API_WORKERS: ${API_WORKERS:-4}
|
||||
PM2_ENABLED: ${PM2_ENABLED:-true}
|
||||
STRATUM_WORKERS: ${STRATUM_WORKERS:-2}
|
||||
STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330}
|
||||
|
||||
@@ -79,6 +79,7 @@ services:
|
||||
REDIS_URL: redis://redis:6379
|
||||
API_PORT: ${API_PORT:-3334}
|
||||
API_SECURE: ${API_SECURE:-false}
|
||||
API_WORKERS: ${API_WORKERS:-4}
|
||||
PM2_ENABLED: ${PM2_ENABLED:-true}
|
||||
STRATUM_WORKERS: ${STRATUM_WORKERS:-2}
|
||||
STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330}
|
||||
|
||||
+2
-2
@@ -4,8 +4,8 @@ module.exports = {
|
||||
{
|
||||
name: 'api',
|
||||
script: './dist/main.js',
|
||||
instances: 1,
|
||||
exec_mode: 'fork',
|
||||
instances: parseInt(process.env.API_WORKERS || '4', 10),
|
||||
exec_mode: 'cluster',
|
||||
env: {
|
||||
MASTER: 'false',
|
||||
API_ONLY: 'true',
|
||||
|
||||
@@ -97,6 +97,7 @@ services:
|
||||
DB_PASSWORD: public_pool
|
||||
DB_DATABASE: public_pool_mainnet
|
||||
REDIS_URL: redis://redis:6379
|
||||
API_WORKERS: ${API_WORKERS:-4}
|
||||
PM2_ENABLED: "true"
|
||||
STRATUM_WORKERS: ${STRATUM_WORKERS:-2}
|
||||
|
||||
|
||||
@@ -0,0 +1,32 @@
|
||||
import { MigrationInterface, QueryRunner } from 'typeorm';
|
||||
|
||||
export class AcceptedShareRollupIndexes1780865400000 implements MigrationInterface {
|
||||
public name = 'AcceptedShareRollupIndexes1780865400000';
|
||||
public transaction = false;
|
||||
|
||||
public async up(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_bucket"
|
||||
ON "accepted_share_10m" ("bucket" DESC)
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_address_bucket"
|
||||
ON "accepted_share_10m" ("address", "bucket" DESC)
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_group_bucket"
|
||||
ON "accepted_share_10m" ("address", "clientName", "bucket" DESC)
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_client_bucket"
|
||||
ON "accepted_share_10m" ("clientId", "bucket" DESC)
|
||||
`);
|
||||
}
|
||||
|
||||
public async down(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_client_bucket"`);
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_group_bucket"`);
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_address_bucket"`);
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_bucket"`);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
import { MigrationInterface, QueryRunner } from 'typeorm';
|
||||
|
||||
export class CurrentRoundBestShareIndex1780897600000 implements MigrationInterface {
|
||||
public name = 'CurrentRoundBestShareIndex1780897600000';
|
||||
public transaction = false;
|
||||
|
||||
public async up(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_address_settings_best_difficulty"`);
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_round_best"
|
||||
ON "accepted_share_entity" ("submissionDifficulty" DESC, "acceptedAt" DESC, "blockHeight")
|
||||
`);
|
||||
}
|
||||
|
||||
public async down(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_round_best"`);
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_address_settings_best_difficulty"
|
||||
ON "address_settings_entity" ("bestDifficulty" DESC)
|
||||
`);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
import { MigrationInterface, QueryRunner } from 'typeorm';
|
||||
|
||||
export class PoolAccountingDashboardIndexes1780867200000 implements MigrationInterface {
|
||||
public name = 'PoolAccountingDashboardIndexes1780867200000';
|
||||
public transaction = false;
|
||||
|
||||
public async up(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_accepted_at"
|
||||
ON "accepted_share_entity" ("acceptedAt" DESC)
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_round_best"
|
||||
ON "accepted_share_entity" ("submissionDifficulty" DESC, "acceptedAt" DESC, "blockHeight")
|
||||
`);
|
||||
}
|
||||
|
||||
public async down(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_round_best"`);
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_accepted_at"`);
|
||||
}
|
||||
}
|
||||
@@ -5,6 +5,7 @@ import { Repository } from 'typeorm';
|
||||
import { ClientEntity } from '../../client/client.entity';
|
||||
import { UserAgentReportView } from './user-agent-report.view';
|
||||
import { RedisMessagingService } from '../../../services/redis-messaging.service';
|
||||
import { logTiming, timeAsync, timingStart } from '../../../utils/timing.utils';
|
||||
|
||||
@Injectable()
|
||||
export class UserAgentReportService {
|
||||
@@ -22,21 +23,27 @@ export class UserAgentReportService {
|
||||
}
|
||||
|
||||
public async getReport() {
|
||||
const cachedReport = await this.redisMessagingService
|
||||
const start = timingStart();
|
||||
const cachedReport = await timeAsync('user agent report shared cache read', () => this.redisMessagingService
|
||||
.getJsonCache<UserAgentReportView[]>(this.liveReportCacheKey)
|
||||
.catch(error => {
|
||||
console.error(`Live user-agent report cache read failed: ${error.message}`);
|
||||
return null;
|
||||
});
|
||||
}));
|
||||
if (cachedReport != null) {
|
||||
logTiming('user agent report getReport', start, { cache: 'hit', rows: cachedReport.length });
|
||||
return cachedReport;
|
||||
}
|
||||
|
||||
if (process.env.API_ONLY == 'true') {
|
||||
return this.userAgentReport.find();
|
||||
const rows = await timeAsync('user agent report materialized view read', () => this.userAgentReport.find());
|
||||
logTiming('user agent report getReport', start, { cache: 'miss', source: 'materialized-view', rows: rows.length });
|
||||
return rows;
|
||||
}
|
||||
|
||||
return this.refreshLiveReport();
|
||||
const rows = await this.refreshLiveReport();
|
||||
logTiming('user agent report getReport', start, { cache: 'miss', source: 'live-presence', rows: rows.length });
|
||||
return rows;
|
||||
}
|
||||
|
||||
public async refreshLiveReport() {
|
||||
@@ -53,7 +60,8 @@ export class UserAgentReportService {
|
||||
}
|
||||
|
||||
private async buildLiveReport() {
|
||||
const presences = await this.redisMessagingService.getAllClientPresence();
|
||||
const start = timingStart();
|
||||
const presences = await timeAsync('user agent report all presence load', () => this.redisMessagingService.getAllClientPresence());
|
||||
const rows = new Map<string, {
|
||||
userAgent: string;
|
||||
count: number;
|
||||
@@ -95,6 +103,7 @@ export class UserAgentReportService {
|
||||
console.error(`Live user-agent report cache write failed: ${error.message}`);
|
||||
});
|
||||
|
||||
logTiming('user agent report buildLiveReport', start, { presences: presences.length, rows: report.length });
|
||||
return report;
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import { Injectable } from '@nestjs/common';
|
||||
import { InjectDataSource } from '@nestjs/typeorm';
|
||||
import { DataSource } from 'typeorm';
|
||||
import { timeAsync } from '../../utils/timing.utils';
|
||||
|
||||
const HASHES_PER_DIFFICULTY = 4294967296;
|
||||
const CHART_BUCKET_SECONDS = 600;
|
||||
@@ -31,14 +32,14 @@ export class ClientStatisticsService {
|
||||
}
|
||||
|
||||
public async getHashRateForGroup(address: string, clientName: string) {
|
||||
const result = await this.dataSource.query(`
|
||||
const result = await timeAsync('client statistics group hashrate query', () => this.dataSource.query(`
|
||||
SELECT
|
||||
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]);
|
||||
`, [address, clientName]), { address, clientName });
|
||||
|
||||
return parseFloat(result[0]?.hashRate ?? '0');
|
||||
}
|
||||
@@ -90,7 +91,12 @@ export class ClientStatisticsService {
|
||||
ORDER BY "label"
|
||||
`;
|
||||
|
||||
const result = await this.dataSource.query(query, params);
|
||||
const result = await timeAsync('client statistics chart query', () => this.dataSource.query(query, params), {
|
||||
filterSql,
|
||||
params: params.length,
|
||||
limit,
|
||||
windowSql,
|
||||
});
|
||||
|
||||
return result.map(res => {
|
||||
return {
|
||||
|
||||
@@ -121,7 +121,7 @@ describe('ShareAccountingService', () => {
|
||||
await expect(queued).resolves.toEqual(expect.objectContaining({ jobId: 'queued' }));
|
||||
});
|
||||
|
||||
it('should return numeric accounting summaries with protocol breakdown', async () => {
|
||||
it('should return numeric accounting summaries from the share rollup', async () => {
|
||||
const repository = {
|
||||
query: jest.fn()
|
||||
.mockResolvedValueOnce([{
|
||||
@@ -135,14 +135,7 @@ describe('ShareAccountingService', () => {
|
||||
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);
|
||||
@@ -158,24 +151,144 @@ describe('ShareAccountingService', () => {
|
||||
creditedDifficultyLastDay: 96,
|
||||
hashRateLast10Minutes: 458129844.9,
|
||||
hashRateLastHour: 114532461.2,
|
||||
bestSubmissionDifficulty: 2048,
|
||||
blockCandidateCount: 1,
|
||||
bestSubmissionDifficulty: 0,
|
||||
bestSubmissionDifficultyAt: null,
|
||||
blockCandidateCount: 0,
|
||||
latestShareAt: '2026-06-07T12:10:00.000Z',
|
||||
protocolBreakdown: [{
|
||||
protocol: 'sv1',
|
||||
acceptedShares: 3,
|
||||
creditedDifficulty: 96,
|
||||
}],
|
||||
protocolBreakdown: [],
|
||||
});
|
||||
expect(repository.query).toHaveBeenNthCalledWith(
|
||||
1,
|
||||
expect.stringContaining('"address" = $1'),
|
||||
expect.stringContaining('"accepted_share_10m"'),
|
||||
['bc1qtest'],
|
||||
);
|
||||
expect(repository.query).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it('should serve API-only address accounting from the share rollup instead of empty summaries', async () => {
|
||||
process.env.API_ONLY = 'true';
|
||||
const repository = {
|
||||
query: jest.fn()
|
||||
.mockResolvedValueOnce([{
|
||||
totalAcceptedShares: '12',
|
||||
totalCreditedDifficulty: '384',
|
||||
acceptedSharesLast10Minutes: '4',
|
||||
creditedDifficultyLast10Minutes: '128',
|
||||
acceptedSharesLastHour: '10',
|
||||
creditedDifficultyLastHour: '320',
|
||||
acceptedSharesLastDay: '12',
|
||||
creditedDifficultyLastDay: '384',
|
||||
hashRateLast10Minutes: '916259689.8',
|
||||
hashRateLastHour: '381774870.2',
|
||||
latestShareAt: new Date('2026-06-07T12:30:00Z'),
|
||||
}]),
|
||||
};
|
||||
const service = new ShareAccountingService(repository as any);
|
||||
|
||||
await expect(service.getAddressSummary('bc1qapi')).resolves.toEqual(expect.objectContaining({
|
||||
totalAcceptedShares: 12,
|
||||
totalCreditedDifficulty: 384,
|
||||
acceptedSharesLast10Minutes: 4,
|
||||
creditedDifficultyLast10Minutes: 128,
|
||||
hashRateLast10Minutes: 916259689.8,
|
||||
latestShareAt: '2026-06-07T12:30:00.000Z',
|
||||
}));
|
||||
expect(repository.query).toHaveBeenCalledWith(
|
||||
expect.stringContaining('FROM "accepted_share_10m"'),
|
||||
['bc1qapi'],
|
||||
);
|
||||
expect(repository.query).not.toHaveBeenCalledWith(
|
||||
expect.stringContaining('FROM "accepted_share_entity"'),
|
||||
expect.anything(),
|
||||
);
|
||||
});
|
||||
|
||||
it('should use Redis cache for share accounting summaries across API workers', async () => {
|
||||
process.env.SHARE_ACCOUNTING_SUMMARY_CACHE_MS = '0';
|
||||
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',
|
||||
latestShareAt: null,
|
||||
}]),
|
||||
};
|
||||
const redis = {
|
||||
getJsonCache: jest.fn()
|
||||
.mockResolvedValueOnce(null)
|
||||
.mockResolvedValueOnce({
|
||||
...new ShareAccountingService(repository as any).emptySummary(),
|
||||
totalAcceptedShares: 1,
|
||||
}),
|
||||
setJsonCache: jest.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
const service = new ShareAccountingService(repository as any, redis as any);
|
||||
|
||||
await service.getAddressSummary('bc1qcached');
|
||||
await expect(service.getAddressSummary('bc1qcached')).resolves.toEqual(expect.objectContaining({
|
||||
totalAcceptedShares: 1,
|
||||
}));
|
||||
|
||||
expect(repository.query).toHaveBeenCalledTimes(1);
|
||||
expect(redis.setJsonCache).toHaveBeenCalledWith(
|
||||
expect.stringContaining('accounting:summary:'),
|
||||
expect.objectContaining({ totalAcceptedShares: 1 }),
|
||||
30000,
|
||||
);
|
||||
});
|
||||
|
||||
it('should overlay live pool data and best share from the current round', async () => {
|
||||
const redis = {
|
||||
getJsonCache: jest.fn().mockResolvedValue(null),
|
||||
setJsonCache: jest.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
const repository = {
|
||||
query: jest.fn()
|
||||
.mockResolvedValueOnce([{
|
||||
totalAcceptedShares: '3',
|
||||
totalCreditedDifficulty: '96',
|
||||
acceptedSharesLast10Minutes: '0',
|
||||
creditedDifficultyLast10Minutes: '0',
|
||||
acceptedSharesLastHour: '3',
|
||||
creditedDifficultyLastHour: '96',
|
||||
acceptedSharesLastDay: '3',
|
||||
creditedDifficultyLastDay: '96',
|
||||
hashRateLast10Minutes: '0',
|
||||
hashRateLastHour: '114532461.2',
|
||||
latestShareAt: new Date('2026-06-07T12:10:00Z'),
|
||||
}])
|
||||
.mockResolvedValueOnce([{
|
||||
acceptedSharesLast10Minutes: '7',
|
||||
creditedDifficultyLast10Minutes: '224',
|
||||
hashRateLast10Minutes: '1603451170.77',
|
||||
latestShareAt: new Date('2026-06-07T12:20:00Z'),
|
||||
}])
|
||||
.mockResolvedValueOnce([{
|
||||
bestSubmissionDifficulty: '4096',
|
||||
bestSubmissionDifficultyAt: new Date('2026-06-07T12:19:00Z'),
|
||||
}]),
|
||||
};
|
||||
const service = new ShareAccountingService(repository as any, redis as any);
|
||||
|
||||
await expect(service.refreshPoolSummary()).resolves.toEqual(expect.objectContaining({
|
||||
acceptedSharesLast10Minutes: 7,
|
||||
creditedDifficultyLast10Minutes: 224,
|
||||
hashRateLast10Minutes: 1603451170.77,
|
||||
bestSubmissionDifficulty: 4096,
|
||||
bestSubmissionDifficultyAt: '2026-06-07T12:19:00.000Z',
|
||||
latestShareAt: '2026-06-07T12:20:00.000Z',
|
||||
}));
|
||||
expect(repository.query).toHaveBeenNthCalledWith(
|
||||
2,
|
||||
expect.stringContaining('GROUP BY "protocol"'),
|
||||
['bc1qtest'],
|
||||
3,
|
||||
expect.stringContaining('WHERE "blockHeight" > latest_found_block."height"'),
|
||||
);
|
||||
});
|
||||
|
||||
@@ -194,18 +307,15 @@ describe('ShareAccountingService', () => {
|
||||
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);
|
||||
expect(repository.query).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
@@ -4,6 +4,7 @@ import { Repository } from 'typeorm';
|
||||
|
||||
import { AcceptedShareEntity } from '../accepted-share/accepted-share.entity';
|
||||
import { RedisMessagingService } from '../../services/redis-messaging.service';
|
||||
import { timeAsync } from '../../utils/timing.utils';
|
||||
|
||||
export interface AcceptedShareRecord {
|
||||
protocol: 'sv1' | 'sv2';
|
||||
@@ -38,6 +39,7 @@ export interface ShareAccountingSummary {
|
||||
hashRateLast10Minutes: number;
|
||||
hashRateLastHour: number;
|
||||
bestSubmissionDifficulty: number;
|
||||
bestSubmissionDifficultyAt: string | null;
|
||||
blockCandidateCount: number;
|
||||
latestShareAt: string | null;
|
||||
protocolBreakdown: {
|
||||
@@ -61,11 +63,13 @@ interface AccountingFilter {
|
||||
}
|
||||
|
||||
const HASHES_PER_DIFFICULTY = 4294967296;
|
||||
const ROLLUP_BUCKET_SECONDS = 600;
|
||||
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;
|
||||
const DEFAULT_REDIS_SUMMARY_CACHE_MS = 30000;
|
||||
|
||||
@Injectable()
|
||||
export class ShareAccountingService implements OnModuleDestroy {
|
||||
@@ -79,6 +83,7 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
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);
|
||||
private readonly redisSummaryCacheMs = this.readNonNegativeInt('SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS', DEFAULT_REDIS_SUMMARY_CACHE_MS);
|
||||
|
||||
constructor(
|
||||
@InjectRepository(AcceptedShareEntity)
|
||||
@@ -148,7 +153,7 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
}
|
||||
|
||||
public async refreshPoolSummary(): Promise<ShareAccountingSummary> {
|
||||
const summary = await this.getSummary({});
|
||||
const summary = await this.withPoolLiveOverlay(await this.getSummary({}));
|
||||
await this.redisMessagingService
|
||||
?.setJsonCache(this.poolSummaryCacheKey, summary, 10 * 60 * 1000)
|
||||
.catch(error => {
|
||||
@@ -157,7 +162,52 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
return summary;
|
||||
}
|
||||
|
||||
private emptySummary(): ShareAccountingSummary {
|
||||
private async withPoolLiveOverlay(summary: ShareAccountingSummary): Promise<ShareAccountingSummary> {
|
||||
if (process.env.API_ONLY === 'true') {
|
||||
return summary;
|
||||
}
|
||||
|
||||
const [[liveWindow], [bestDifficultyRow]] = await Promise.all([
|
||||
timeAsync('share accounting live 10m pool query', () => this.acceptedShareRepository.query(`
|
||||
SELECT
|
||||
COUNT(*)::int AS "acceptedSharesLast10Minutes",
|
||||
COALESCE(SUM("creditedDifficulty"), 0)::float AS "creditedDifficultyLast10Minutes",
|
||||
COALESCE((SUM("creditedDifficulty") * ${HASHES_PER_DIFFICULTY}) / 600, 0)::float AS "hashRateLast10Minutes",
|
||||
MAX("acceptedAt") AS "latestShareAt"
|
||||
FROM "accepted_share_entity"
|
||||
WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes'
|
||||
`)),
|
||||
timeAsync('share accounting current round best share query', () => this.acceptedShareRepository.query(`
|
||||
WITH latest_found_block AS (
|
||||
SELECT COALESCE(MAX("height"), 0) AS "height"
|
||||
FROM "blocks_entity"
|
||||
)
|
||||
SELECT
|
||||
COALESCE("submissionDifficulty", 0)::float AS "bestSubmissionDifficulty",
|
||||
"acceptedAt" AS "bestSubmissionDifficultyAt"
|
||||
FROM "accepted_share_entity", latest_found_block
|
||||
WHERE "blockHeight" > latest_found_block."height"
|
||||
ORDER BY "submissionDifficulty" DESC, "acceptedAt" DESC
|
||||
LIMIT 1
|
||||
`)),
|
||||
]);
|
||||
|
||||
return {
|
||||
...summary,
|
||||
acceptedSharesLast10Minutes: this.toNumber(liveWindow?.acceptedSharesLast10Minutes),
|
||||
creditedDifficultyLast10Minutes: this.toNumber(liveWindow?.creditedDifficultyLast10Minutes),
|
||||
hashRateLast10Minutes: this.toNumber(liveWindow?.hashRateLast10Minutes),
|
||||
bestSubmissionDifficulty: this.toNumber(bestDifficultyRow?.bestSubmissionDifficulty),
|
||||
bestSubmissionDifficultyAt: bestDifficultyRow?.bestSubmissionDifficultyAt == null
|
||||
? null
|
||||
: new Date(bestDifficultyRow.bestSubmissionDifficultyAt).toISOString(),
|
||||
latestShareAt: liveWindow?.latestShareAt == null
|
||||
? summary.latestShareAt
|
||||
: new Date(liveWindow.latestShareAt).toISOString(),
|
||||
};
|
||||
}
|
||||
|
||||
public emptySummary(): ShareAccountingSummary {
|
||||
return {
|
||||
totalAcceptedShares: 0,
|
||||
totalCreditedDifficulty: 0,
|
||||
@@ -170,6 +220,7 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
hashRateLast10Minutes: 0,
|
||||
hashRateLastHour: 0,
|
||||
bestSubmissionDifficulty: 0,
|
||||
bestSubmissionDifficultyAt: null,
|
||||
blockCandidateCount: 0,
|
||||
latestShareAt: null,
|
||||
protocolBreakdown: [],
|
||||
@@ -191,20 +242,19 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
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) {
|
||||
if (uniqueClientIds.length === 0 || process.env.API_ONLY === 'true') {
|
||||
return summaries;
|
||||
}
|
||||
|
||||
const rows = await this.acceptedShareRepository.query(`
|
||||
const rows = await timeAsync('share accounting session summaries query', () => 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"
|
||||
MAX("bucket") AS "latestShareAt",
|
||||
COALESCE((SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes') * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes"
|
||||
FROM "accepted_share_10m"
|
||||
WHERE "clientId" = ANY($1::uuid[])
|
||||
GROUP BY "clientId"
|
||||
`, [uniqueClientIds]);
|
||||
`, [uniqueClientIds]), { clientIds: uniqueClientIds.length });
|
||||
|
||||
rows.forEach(row => {
|
||||
summaries.set(row.clientId, {
|
||||
@@ -213,7 +263,7 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
? null
|
||||
: new Date(row.latestShareAt).toISOString(),
|
||||
hashRateLast10Minutes: this.toNumber(row.hashRateLast10Minutes),
|
||||
bestSubmissionDifficulty: this.toNumber(row.bestSubmissionDifficulty),
|
||||
bestSubmissionDifficulty: 0,
|
||||
});
|
||||
});
|
||||
|
||||
@@ -229,7 +279,7 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
return cached.value;
|
||||
}
|
||||
|
||||
const value = this.loadSummary(filter).catch(error => {
|
||||
const value = this.loadCachedSummary(filter, cacheKey).catch(error => {
|
||||
this.summaryCache.delete(cacheKey);
|
||||
throw error;
|
||||
});
|
||||
@@ -245,37 +295,47 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
return value;
|
||||
}
|
||||
|
||||
private async loadCachedSummary(filter: AccountingFilter, cacheKey: string): Promise<ShareAccountingSummary> {
|
||||
const redisCacheKey = `accounting:summary:${cacheKey}`;
|
||||
const cached = await this.redisMessagingService
|
||||
?.getJsonCache<ShareAccountingSummary>(redisCacheKey)
|
||||
.catch(error => {
|
||||
console.error(`Share accounting summary cache read failed: ${error.message}`);
|
||||
return null;
|
||||
});
|
||||
if (cached != null) {
|
||||
return cached;
|
||||
}
|
||||
|
||||
const summary = await this.loadSummary(filter);
|
||||
|
||||
await this.redisMessagingService
|
||||
?.setJsonCache(redisCacheKey, summary, this.redisSummaryCacheMs)
|
||||
.catch(error => {
|
||||
console.error(`Share accounting summary cache write failed: ${error.message}`);
|
||||
});
|
||||
|
||||
return summary;
|
||||
}
|
||||
|
||||
private async loadSummary(filter: AccountingFilter): Promise<ShareAccountingSummary> {
|
||||
const { whereSql, params } = this.buildWhereClause(filter);
|
||||
const [summary] = await this.acceptedShareRepository.query(`
|
||||
const [summary] = await timeAsync('share accounting summary query', () => 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"
|
||||
COALESCE(SUM("acceptedCount"), 0)::int AS "totalAcceptedShares",
|
||||
COALESCE(SUM("shares"), 0)::float AS "totalCreditedDifficulty",
|
||||
COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes'), 0)::int AS "acceptedSharesLast10Minutes",
|
||||
COALESCE(SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes'), 0)::float AS "creditedDifficultyLast10Minutes",
|
||||
COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 hour'), 0)::int AS "acceptedSharesLastHour",
|
||||
COALESCE(SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 hour'), 0)::float AS "creditedDifficultyLastHour",
|
||||
COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 day'), 0)::int AS "acceptedSharesLastDay",
|
||||
COALESCE(SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 day'), 0)::float AS "creditedDifficultyLastDay",
|
||||
COALESCE((SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '10 minutes') * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes",
|
||||
COALESCE((SUM("shares") FILTER (WHERE "bucket" > NOW() - INTERVAL '1 hour') * ${HASHES_PER_DIFFICULTY}) / 3600, 0)::float AS "hashRateLastHour",
|
||||
MAX("bucket") AS "latestShareAt"
|
||||
FROM "accepted_share_10m"
|
||||
${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);
|
||||
`, params), { filter, params: params.length });
|
||||
|
||||
return {
|
||||
totalAcceptedShares: this.toNumber(summary?.totalAcceptedShares),
|
||||
@@ -288,16 +348,13 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
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),
|
||||
bestSubmissionDifficulty: 0,
|
||||
bestSubmissionDifficultyAt: null,
|
||||
blockCandidateCount: 0,
|
||||
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),
|
||||
})),
|
||||
protocolBreakdown: [],
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
+107
-32
@@ -13,11 +13,13 @@ import { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-r
|
||||
import { StratumV2Service } from './services/stratum-v2.service';
|
||||
import { ShareAccountingService } from './ORM/share-accounting/share-accounting.service';
|
||||
import { RedisMessagingService } from './services/redis-messaging.service';
|
||||
import { logTiming, timeAsync, timingStart } from './utils/timing.utils';
|
||||
|
||||
@Controller()
|
||||
export class AppController {
|
||||
|
||||
private uptime = new Date();
|
||||
private siteInfoRefreshPromise: Promise<SiteInfoResponse> | null = null;
|
||||
|
||||
constructor(
|
||||
@Inject(CACHE_MANAGER) private readonly cacheManager: Cache,
|
||||
@@ -34,30 +36,59 @@ export class AppController {
|
||||
|
||||
@Get('info')
|
||||
public async info() {
|
||||
const start = timingStart();
|
||||
|
||||
|
||||
const CACHE_KEY = 'SITE_INFO';
|
||||
const cachedResult = await this.getCached(CACHE_KEY);
|
||||
const STALE_CACHE_KEY = 'SITE_INFO_STALE';
|
||||
const cachedResult = await this.getCached(CACHE_KEY, 5 * 60 * 1000);
|
||||
|
||||
if (cachedResult != null) {
|
||||
logTiming('GET /api/info', start, { cache: 'fresh' });
|
||||
return cachedResult;
|
||||
}
|
||||
|
||||
let usedFallback = false;
|
||||
const withInfoTimeout = async <T>(label: string, promise: Promise<T>, fallback: T): Promise<T> => {
|
||||
return await this.withTimeout(label, promise, fallback, () => {
|
||||
usedFallback = true;
|
||||
const staleResult = await this.getCached<SiteInfoResponse>(STALE_CACHE_KEY, 60 * 60 * 1000);
|
||||
if (staleResult != null) {
|
||||
void this.refreshSiteInfo(staleResult);
|
||||
logTiming('GET /api/info', start, { cache: 'stale' });
|
||||
return staleResult;
|
||||
}
|
||||
|
||||
const response = await this.refreshSiteInfo(null);
|
||||
logTiming('GET /api/info', start, { cache: 'miss' });
|
||||
return response;
|
||||
|
||||
}
|
||||
|
||||
private async refreshSiteInfo(staleInfo: SiteInfoResponse | null): Promise<SiteInfoResponse> {
|
||||
if (this.siteInfoRefreshPromise != null) {
|
||||
return this.siteInfoRefreshPromise;
|
||||
}
|
||||
|
||||
this.siteInfoRefreshPromise = this.loadSiteInfo(staleInfo)
|
||||
.finally(() => {
|
||||
this.siteInfoRefreshPromise = null;
|
||||
});
|
||||
|
||||
return this.siteInfoRefreshPromise;
|
||||
}
|
||||
|
||||
private async loadSiteInfo(staleInfo: SiteInfoResponse | null): Promise<SiteInfoResponse> {
|
||||
const CACHE_KEY = 'SITE_INFO';
|
||||
const STALE_CACHE_KEY = 'SITE_INFO_STALE';
|
||||
const withInfoTimeout = async <T>(label: string, promise: Promise<T>, fallback: T): Promise<T> => {
|
||||
return await this.withTimeout(label, promise, fallback);
|
||||
};
|
||||
|
||||
const [blockData, highScores, poolAuthority, userAgentReport] = await Promise.all([
|
||||
withInfoTimeout('found blocks', this.blocksService.getFoundBlocks(), []),
|
||||
withInfoTimeout('high scores', this.addressSettingsService.getHighScores(), []),
|
||||
withInfoTimeout('SV2 authority', this.stratumV2Service.getPoolAuthorityPublicKey(), {
|
||||
publicKey: '',
|
||||
configured: false
|
||||
withInfoTimeout('found blocks', timeAsync('/api/info found blocks', () => this.blocksService.getFoundBlocks()), staleInfo?.blockData ?? []),
|
||||
withInfoTimeout('high scores', timeAsync('/api/info high scores', () => this.addressSettingsService.getHighScores()), staleInfo?.highScores ?? []),
|
||||
withInfoTimeout('SV2 authority', timeAsync('/api/info SV2 authority', () => this.stratumV2Service.getPoolAuthorityPublicKey()), {
|
||||
publicKey: staleInfo?.sv2?.poolAuthorityPublicKey ?? '',
|
||||
configured: staleInfo?.sv2?.authorityKeyConfigured ?? false
|
||||
}),
|
||||
withInfoTimeout<UserAgentReportView[]>('user agent report', this.userAgentReportService.getReport(), []),
|
||||
withInfoTimeout<UserAgentReportView[]>('user agent report', timeAsync('/api/info user agent report', () => this.userAgentReportService.getReport()), staleInfo?.userAgents ?? []),
|
||||
]);
|
||||
|
||||
const other: {
|
||||
@@ -87,7 +118,7 @@ export class AppController {
|
||||
userAgents.push({ userAgent: 'Other', count: other.count.toString(), bestDifficulty: other.bestDifficulty, totalHashRate: other.totalHashRate.toString() })
|
||||
}
|
||||
|
||||
const data = {
|
||||
const data: SiteInfoResponse = {
|
||||
blockData,
|
||||
userAgents,
|
||||
highScores,
|
||||
@@ -98,46 +129,52 @@ export class AppController {
|
||||
uptime: this.uptime
|
||||
};
|
||||
|
||||
// Match the pre-Timescale dashboard cache behavior; live accounting is exposed separately.
|
||||
await this.setCached(CACHE_KEY, data, usedFallback ? 15 * 1000 : 5 * 60 * 1000);
|
||||
// Cache a complete response even when one slow component falls back to stale data.
|
||||
// A short retry loop here causes repeated DB work because timed-out TypeORM queries are not cancelled.
|
||||
await this.setCached(CACHE_KEY, data, 5 * 60 * 1000);
|
||||
await this.setCached(STALE_CACHE_KEY, data, 60 * 60 * 1000);
|
||||
|
||||
return data;
|
||||
|
||||
}
|
||||
|
||||
@Get('info/accounting')
|
||||
public async infoAccounting() {
|
||||
const start = timingStart();
|
||||
const CACHE_KEY = 'SITE_ACCOUNTING';
|
||||
const cachedResult = await this.getCached(CACHE_KEY);
|
||||
const cachedResult = await this.getCached(CACHE_KEY, 15 * 1000);
|
||||
|
||||
if (cachedResult != null) {
|
||||
logTiming('GET /api/info/accounting', start, { cache: 'hit' });
|
||||
return cachedResult;
|
||||
}
|
||||
|
||||
const data = await this.shareAccountingService.getPoolSummary();
|
||||
const data = await timeAsync('/api/info/accounting getPoolSummary', () => this.shareAccountingService.getPoolSummary());
|
||||
|
||||
//15 sec
|
||||
await this.setCached(CACHE_KEY, data, 15 * 1000);
|
||||
|
||||
logTiming('GET /api/info/accounting', start, { cache: 'miss' });
|
||||
return data;
|
||||
}
|
||||
|
||||
@Get('pool')
|
||||
public async pool() {
|
||||
const start = timingStart();
|
||||
|
||||
const CACHE_KEY = 'POOL_INFO';
|
||||
const cachedResult = await this.getCached(CACHE_KEY);
|
||||
const cachedResult = await this.getCached(CACHE_KEY, 15 * 1000);
|
||||
|
||||
if (cachedResult != null) {
|
||||
logTiming('GET /api/pool', start, { cache: 'hit' });
|
||||
return cachedResult;
|
||||
}
|
||||
|
||||
const userAgents = await this.userAgentReportService.getReport();
|
||||
const userAgents = await timeAsync('/api/pool user agent report', () => this.userAgentReportService.getReport());
|
||||
|
||||
const totalHashRate = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.totalHashRate), 0);
|
||||
const totalMiners = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.count), 0);
|
||||
const blockHeight = this.bitcoinRpcService.miningInfo.blocks;
|
||||
const blocksFound = await this.blocksService.getFoundBlocks();
|
||||
const blocksFound = await timeAsync('/api/pool found blocks', () => this.blocksService.getFoundBlocks());
|
||||
|
||||
const data = {
|
||||
totalHashRate,
|
||||
@@ -150,61 +187,74 @@ export class AppController {
|
||||
// Keep online miner counts responsive after reconnect cleanup.
|
||||
await this.setCached(CACHE_KEY, data, 15 * 1000);
|
||||
|
||||
logTiming('GET /api/pool', start, { cache: 'miss', userAgentCount: userAgents.length });
|
||||
return data;
|
||||
}
|
||||
|
||||
@Get('network')
|
||||
public async network() {
|
||||
const start = timingStart();
|
||||
logTiming('GET /api/network', start);
|
||||
return this.bitcoinRpcService.miningInfo ?? {};
|
||||
}
|
||||
|
||||
@Get('info/chart')
|
||||
public async infoChart() {
|
||||
const start = timingStart();
|
||||
|
||||
|
||||
const CACHE_KEY = 'SITE_HASHRATE_GRAPH';
|
||||
const cachedResult = await this.getCached(CACHE_KEY);
|
||||
const cachedResult = await this.getCached(CACHE_KEY, 10 * 60 * 1000);
|
||||
|
||||
if (cachedResult != null) {
|
||||
logTiming('GET /api/info/chart', start, { cache: 'hit' });
|
||||
return cachedResult;
|
||||
}
|
||||
|
||||
const chartData = await this.clientStatisticsService.getChartDataForSite();
|
||||
const chartData = await timeAsync('/api/info/chart query', () => this.clientStatisticsService.getChartDataForSite());
|
||||
|
||||
//10 min
|
||||
await this.setCached(CACHE_KEY, chartData, 10 * 60 * 1000);
|
||||
|
||||
logTiming('GET /api/info/chart', start, { cache: 'miss', points: chartData.length });
|
||||
return chartData;
|
||||
|
||||
|
||||
}
|
||||
|
||||
private async getCached<T>(key: string): Promise<T | null> {
|
||||
const shared = await this.redisMessagingService.getJsonCache<T>(`api:${key}`).catch(error => {
|
||||
private async getCached<T>(key: string, localTtlMs: number): Promise<T | null> {
|
||||
const local = await this.cacheManager.get<T>(key);
|
||||
if (local != null) {
|
||||
return local;
|
||||
}
|
||||
|
||||
const shared = await this.withSecondaryCacheTimeout(
|
||||
this.redisMessagingService.getJsonCache<T>(`api:${key}`).catch(error => {
|
||||
console.error(`Shared API cache read failed for ${key}: ${error.message}`);
|
||||
return null;
|
||||
});
|
||||
}),
|
||||
100
|
||||
);
|
||||
if (shared != null) {
|
||||
await this.cacheManager.set(key, shared, localTtlMs);
|
||||
return shared;
|
||||
}
|
||||
|
||||
return await this.cacheManager.get<T>(key) ?? null;
|
||||
return null;
|
||||
}
|
||||
|
||||
private async setCached(key: string, value: unknown, ttlMs: number): Promise<void> {
|
||||
await Promise.all([
|
||||
this.cacheManager.set(key, value, ttlMs),
|
||||
this.redisMessagingService.setJsonCache(`api:${key}`, value, ttlMs).catch(error => {
|
||||
await this.cacheManager.set(key, value, ttlMs);
|
||||
void this.redisMessagingService.setJsonCache(`api:${key}`, value, ttlMs).catch(error => {
|
||||
console.error(`Shared API cache write failed for ${key}: ${error.message}`);
|
||||
}),
|
||||
]);
|
||||
});
|
||||
}
|
||||
|
||||
private async withTimeout<T>(
|
||||
label: string,
|
||||
promise: Promise<T>,
|
||||
fallback: T,
|
||||
onTimeout: () => void,
|
||||
onTimeout: () => void = () => undefined,
|
||||
timeoutMs = 1500
|
||||
): Promise<T> {
|
||||
let timeout: NodeJS.Timeout;
|
||||
@@ -224,4 +274,29 @@ export class AppController {
|
||||
}
|
||||
}
|
||||
|
||||
private async withSecondaryCacheTimeout<T>(promise: Promise<T | null>, timeoutMs: number): Promise<T | null> {
|
||||
let timeout: NodeJS.Timeout;
|
||||
const timeoutPromise = new Promise<null>(resolve => {
|
||||
timeout = setTimeout(() => resolve(null), timeoutMs);
|
||||
timeout.unref?.();
|
||||
});
|
||||
|
||||
try {
|
||||
return await Promise.race([promise, timeoutPromise]);
|
||||
} finally {
|
||||
clearTimeout(timeout);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
interface SiteInfoResponse {
|
||||
blockData: unknown[];
|
||||
userAgents: UserAgentReportView[];
|
||||
highScores: unknown[];
|
||||
sv2: {
|
||||
poolAuthorityPublicKey: string;
|
||||
authorityKeyConfigured: boolean;
|
||||
};
|
||||
uptime: Date;
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import { ClientStatisticsService } from '../../ORM/client-statistics/client-stat
|
||||
import { ClientService } from '../../ORM/client/client.service';
|
||||
import { ShareAccountingService } from '../../ORM/share-accounting/share-accounting.service';
|
||||
import { RedisMessagingService } from '../../services/redis-messaging.service';
|
||||
import { logTiming, timeAsync, timingStart } from '../../utils/timing.utils';
|
||||
|
||||
|
||||
@Controller('client')
|
||||
@@ -21,16 +22,23 @@ export class ClientController {
|
||||
|
||||
@Get(':address')
|
||||
async getClientInfo(@Param('address') address: string) {
|
||||
const start = timingStart();
|
||||
|
||||
const workers = await this.redisMessagingService.getClientPresenceByAddress(address);
|
||||
const sessionSummaries = await this.shareAccountingService.getSessionSummaries(workers.map(worker => worker.clientId));
|
||||
const workers = await timeAsync('/api/client/:address presence', () => this.redisMessagingService.getClientPresenceByAddress(address), { address });
|
||||
const sessionSummaries = await timeAsync('/api/client/:address session summaries', () => this.shareAccountingService.getSessionSummaries(workers.map(worker => worker.clientId)), { address, workers: workers.length });
|
||||
|
||||
const addressSettings = await this.addressSettingsService.getSettings(address, false);
|
||||
const addressSettings = process.env.API_ONLY === 'true'
|
||||
? null
|
||||
: await timeAsync('/api/client/:address address settings', () => this.addressSettingsService.getSettings(address, false), { address });
|
||||
const accounting = await timeAsync('/api/client/:address accounting', () => this.shareAccountingService.getAddressSummary(address), { address });
|
||||
const bestDifficulty = addressSettings?.bestDifficulty ?? workers.reduce((best, worker) => {
|
||||
return Math.max(best, Number(worker.bestDifficulty ?? 0));
|
||||
}, 0);
|
||||
|
||||
return {
|
||||
bestDifficulty: addressSettings?.bestDifficulty,
|
||||
const response = {
|
||||
bestDifficulty,
|
||||
workersCount: workers.length,
|
||||
accounting: await this.shareAccountingService.getAddressSummary(address),
|
||||
accounting,
|
||||
workers: await Promise.all(
|
||||
workers.map(async (worker) => {
|
||||
const sessionSummary = sessionSummaries.get(worker.clientId);
|
||||
@@ -49,18 +57,24 @@ export class ClientController {
|
||||
})
|
||||
)
|
||||
}
|
||||
logTiming('GET /api/client/:address', start, { address, workers: workers.length });
|
||||
return response;
|
||||
}
|
||||
|
||||
@Get(':address/chart')
|
||||
async getClientInfoChart(@Param('address') address: string) {
|
||||
const chartData = await this.clientStatisticsService.getChartDataForAddress(address);
|
||||
const start = timingStart();
|
||||
const chartData = await timeAsync('/api/client/:address/chart query', () => this.clientStatisticsService.getChartDataForAddress(address), { address });
|
||||
logTiming('GET /api/client/:address/chart', start, { address, points: chartData.length });
|
||||
return chartData;
|
||||
}
|
||||
|
||||
@Get(':address/:workerName')
|
||||
async getWorkerGroupInfo(@Param('address') address: string, @Param('workerName') workerName: string) {
|
||||
const start = timingStart();
|
||||
|
||||
const workers = (await this.redisMessagingService.getClientPresenceByAddress(address))
|
||||
const addressWorkers = await timeAsync('/api/client/:address/:workerName address presence', () => this.redisMessagingService.getClientPresenceByAddress(address), { address, workerName });
|
||||
const workers = addressWorkers
|
||||
.filter(worker => worker.clientName === workerName);
|
||||
|
||||
const bestDifficulty = workers.reduce((pre, cur, idx, arr) => {
|
||||
@@ -70,24 +84,29 @@ export class ClientController {
|
||||
return pre;
|
||||
}, 0);
|
||||
|
||||
const chartData = await this.clientStatisticsService.getChartDataForGroup(address, workerName);
|
||||
return {
|
||||
const chartData = await timeAsync('/api/client/:address/:workerName chart', () => this.clientStatisticsService.getChartDataForGroup(address, workerName), { address, workerName });
|
||||
const accounting = await timeAsync('/api/client/:address/:workerName accounting', () => this.shareAccountingService.getWorkerGroupSummary(address, workerName), { address, workerName });
|
||||
const response = {
|
||||
|
||||
name: workerName,
|
||||
bestDifficulty: Math.floor(bestDifficulty),
|
||||
accounting: await this.shareAccountingService.getWorkerGroupSummary(address, workerName),
|
||||
accounting,
|
||||
chartData: chartData,
|
||||
|
||||
}
|
||||
logTiming('GET /api/client/:address/:workerName', start, { address, workerName, addressWorkers: addressWorkers.length, matchedWorkers: workers.length, points: chartData.length });
|
||||
return response;
|
||||
}
|
||||
|
||||
@Get(':address/:workerName/:sessionId')
|
||||
async getWorkerInfo(@Param('address') address: string, @Param('workerName') workerName: string, @Param('sessionId') sessionId: string) {
|
||||
const start = timingStart();
|
||||
|
||||
const presenceWorker = (await this.redisMessagingService.getClientPresenceByAddress(address))
|
||||
const addressWorkers = await timeAsync('/api/client/:address/:workerName/:sessionId address presence', () => this.redisMessagingService.getClientPresenceByAddress(address), { address, workerName, sessionId });
|
||||
const presenceWorker = addressWorkers
|
||||
.find(worker => worker.clientName === workerName && worker.sessionId === sessionId);
|
||||
const worker = presenceWorker == null
|
||||
? await this.clientService.getBySessionId(address, workerName, sessionId)
|
||||
? await timeAsync('/api/client/:address/:workerName/:sessionId DB fallback', () => this.clientService.getBySessionId(address, workerName, sessionId), { address, workerName, sessionId })
|
||||
: {
|
||||
id: presenceWorker.clientId,
|
||||
sessionId: presenceWorker.sessionId,
|
||||
@@ -96,17 +115,21 @@ export class ClientController {
|
||||
startTime: presenceWorker.startTime,
|
||||
};
|
||||
if (worker == null) {
|
||||
logTiming('GET /api/client/:address/:workerName/:sessionId', start, { address, workerName, sessionId, found: false, addressWorkers: addressWorkers.length });
|
||||
return new NotFoundException();
|
||||
}
|
||||
const chartData = await this.clientStatisticsService.getChartDataForSession(worker.id);
|
||||
const chartData = await timeAsync('/api/client/:address/:workerName/:sessionId chart', () => this.clientStatisticsService.getChartDataForSession(worker.id), { address, workerName, sessionId, clientId: worker.id });
|
||||
const accounting = await timeAsync('/api/client/:address/:workerName/:sessionId accounting', () => this.shareAccountingService.getSessionSummary(worker.id), { address, workerName, sessionId, clientId: worker.id });
|
||||
|
||||
return {
|
||||
const response = {
|
||||
sessionId: worker.sessionId,
|
||||
name: worker.clientName,
|
||||
bestDifficulty: Math.floor(worker.bestDifficulty),
|
||||
accounting: await this.shareAccountingService.getSessionSummary(worker.id),
|
||||
accounting,
|
||||
chartData: chartData,
|
||||
startTime: worker.startTime
|
||||
}
|
||||
logTiming('GET /api/client/:address/:workerName/:sessionId', start, { address, workerName, sessionId, found: true, addressWorkers: addressWorkers.length, points: chartData.length });
|
||||
return response;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,6 +5,9 @@ import { InitialTimescaleSchema1780859300000 } from './ORM/_migrations/InitialTi
|
||||
import { ActiveOnlyUserAgentReport1780860200000 } from './ORM/_migrations/ActiveOnlyUserAgentReport1780860200000';
|
||||
import { TimescaleOperationalHardening1780861200000 } from './ORM/_migrations/TimescaleOperationalHardening1780861200000';
|
||||
import { AcceptedShareIndex1780862400000 } from './ORM/_migrations/AcceptedShareIndex1780862400000';
|
||||
import { AcceptedShareRollupIndexes1780865400000 } from './ORM/_migrations/AcceptedShareRollupIndexes1780865400000';
|
||||
import { PoolAccountingDashboardIndexes1780867200000 } from './ORM/_migrations/PoolAccountingDashboardIndexes1780867200000';
|
||||
import { CurrentRoundBestShareIndex1780897600000 } from './ORM/_migrations/CurrentRoundBestShareIndex1780897600000';
|
||||
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';
|
||||
@@ -28,6 +31,9 @@ export const databaseMigrations = [
|
||||
ActiveOnlyUserAgentReport1780860200000,
|
||||
TimescaleOperationalHardening1780861200000,
|
||||
AcceptedShareIndex1780862400000,
|
||||
AcceptedShareRollupIndexes1780865400000,
|
||||
PoolAccountingDashboardIndexes1780867200000,
|
||||
CurrentRoundBestShareIndex1780897600000,
|
||||
];
|
||||
|
||||
export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions {
|
||||
|
||||
@@ -4,6 +4,7 @@ import { createClient, RedisClientType } from 'redis';
|
||||
|
||||
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
|
||||
import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo';
|
||||
import { logTiming, timingStart } from '../utils/timing.utils';
|
||||
|
||||
const MINING_INFO_CHANNEL = 'mining-info.updated';
|
||||
const MINING_INFO_KEY = 'mining-info:latest';
|
||||
@@ -263,8 +264,10 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
|
||||
}
|
||||
|
||||
private async getPresenceFromSet(setKey: string): Promise<ClientPresence[]> {
|
||||
const start = timingStart();
|
||||
const clientIds = await this.publisher.sMembers(setKey);
|
||||
if (clientIds.length === 0) {
|
||||
logTiming('redis presence set load', start, { setKey, clientIds: 0, presences: 0, staleClientIds: 0 });
|
||||
return [];
|
||||
}
|
||||
|
||||
@@ -294,6 +297,12 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
|
||||
}
|
||||
}
|
||||
|
||||
logTiming('redis presence set load', start, {
|
||||
setKey,
|
||||
clientIds: clientIds.length,
|
||||
presences: presences.length,
|
||||
staleClientIds: staleClientIds.length,
|
||||
});
|
||||
return presences;
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,51 @@
|
||||
import { performance } from 'perf_hooks';
|
||||
|
||||
type TimingMetadata = Record<string, unknown> | (() => Record<string, unknown>);
|
||||
|
||||
const DEFAULT_TIMING_LOG_MS = 250;
|
||||
|
||||
export function timingStart(): number {
|
||||
return performance.now();
|
||||
}
|
||||
|
||||
export async function timeAsync<T>(
|
||||
label: string,
|
||||
work: () => Promise<T>,
|
||||
metadata?: TimingMetadata,
|
||||
): Promise<T> {
|
||||
const start = timingStart();
|
||||
try {
|
||||
return await work();
|
||||
} finally {
|
||||
logTiming(label, start, metadata);
|
||||
}
|
||||
}
|
||||
|
||||
export function logTiming(label: string, start: number, metadata?: TimingMetadata): void {
|
||||
const elapsedMs = performance.now() - start;
|
||||
const thresholdMs = getTimingThresholdMs();
|
||||
if (thresholdMs < 0 || elapsedMs < thresholdMs) {
|
||||
return;
|
||||
}
|
||||
|
||||
const resolvedMetadata = resolveMetadata(metadata);
|
||||
const metadataText = resolvedMetadata == null ? '' : ` ${JSON.stringify(resolvedMetadata)}`;
|
||||
console.warn(`[timing] ${label} ${elapsedMs.toFixed(1)}ms${metadataText}`);
|
||||
}
|
||||
|
||||
function getTimingThresholdMs(): number {
|
||||
const configured = Number(process.env.API_TIMING_LOG_MS);
|
||||
return Number.isFinite(configured) ? configured : DEFAULT_TIMING_LOG_MS;
|
||||
}
|
||||
|
||||
function resolveMetadata(metadata?: TimingMetadata): Record<string, unknown> | null {
|
||||
if (metadata == null) {
|
||||
return null;
|
||||
}
|
||||
|
||||
try {
|
||||
return typeof metadata === 'function' ? metadata() : metadata;
|
||||
} catch (error) {
|
||||
return { metadataError: error instanceof Error ? error.message : String(error) };
|
||||
}
|
||||
}
|
||||
@@ -45,6 +45,8 @@ describe('TimescaleDB and Redis integration', () => {
|
||||
await dataSource.query(`DELETE FROM accepted_share_entity`);
|
||||
await dataSource.query(`DELETE FROM client_entity`);
|
||||
await dataSource.query(`REFRESH MATERIALIZED VIEW user_agent_report_view`);
|
||||
await redisMessagingService.clearClientPresence();
|
||||
await (redisMessagingService as any).publisher.del('json-cache:presence:user-agent-report');
|
||||
});
|
||||
|
||||
it('should create Timescale extension, hypertable, continuous aggregates, and operational policies', async () => {
|
||||
|
||||
Reference in New Issue
Block a user