diff --git a/.env.example b/.env.example index baa44c2..69d6bf4 100644 --- a/.env.example +++ b/.env.example @@ -21,15 +21,23 @@ API_BIND_HOST=127.0.0.1 API_PUBLIC_PORT=3334 API_WORKERS=4 +# Docker json-file log rotation. Applies when using the compose files. +DOCKER_LOG_MAX_SIZE=100m +DOCKER_LOG_MAX_FILES=5 + # Plain TCP Stratum ports accept both SV1 JSON-RPC and SV2 Noise/binary traffic. STRATUM_PORTS=3333,3332,3331,3330 STRATUM_WORKERS=2 +STRATUM_MIN_DIFFICULTY=1 +STRATUM_SOCKET_TIMEOUT_MS=3600000 +STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS=60000 # Optional additional SV2-only mining ports. Usually unnecessary when STRATUM_PORTS are exposed. #STRATUM_V2_PORTS= #SV2_START_DIFFICULTY=100000 #SV2_TARGET_SHARES_PER_MINUTE=2 #SV2_DIFFICULTY_CHECK_INTERVAL_MS=60000 +#STRATUM_V2_SOCKET_TIMEOUT_MS=3600000 #SV2_AUTHORITY_PRIVKEY= STRATUM_SECURE=true diff --git a/docker-compose.external-db.yml b/docker-compose.external-db.yml index c18d728..3b3320b 100644 --- a/docker-compose.external-db.yml +++ b/docker-compose.external-db.yml @@ -1,8 +1,15 @@ +x-log-limits: &log-limits + driver: json-file + options: + max-size: ${DOCKER_LOG_MAX_SIZE:-100m} + max-file: "${DOCKER_LOG_MAX_FILES:-5}" + services: redis: image: redis:8-alpine container_name: public-pool-redis restart: unless-stopped + logging: *log-limits ports: - "127.0.0.1:${REDIS_PORT:-6379}:6379/tcp" volumes: @@ -19,6 +26,7 @@ services: context: . dockerfile: Dockerfile restart: unless-stopped + logging: *log-limits depends_on: redis: condition: service_healthy @@ -54,6 +62,9 @@ services: PM2_ENABLED: ${PM2_ENABLED:-true} STRATUM_WORKERS: ${STRATUM_WORKERS:-2} STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330} + STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1} + STRATUM_SOCKET_TIMEOUT_MS: ${STRATUM_SOCKET_TIMEOUT_MS:-3600000} + STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS: ${STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS:-60000} STRATUM_SECURE: ${STRATUM_SECURE:-true} SECURE_STRATUM_PORTS: ${SECURE_STRATUM_PORTS:-4333,4332,4331,4330} STRATUM_MAX_CONNECTIONS_PER_LISTENER: ${STRATUM_MAX_CONNECTIONS_PER_LISTENER:-10000} diff --git a/docker-compose.test.yml b/docker-compose.test.yml index 6add49d..2b5b8f1 100644 --- a/docker-compose.test.yml +++ b/docker-compose.test.yml @@ -1,8 +1,15 @@ name: public-pool-test +x-log-limits: &log-limits + driver: json-file + options: + max-size: ${DOCKER_LOG_MAX_SIZE:-100m} + max-file: "${DOCKER_LOG_MAX_FILES:-5}" + services: timescaledb: image: timescale/timescaledb:latest-pg17 + logging: *log-limits environment: POSTGRES_DB: public_pool_test POSTGRES_USER: public_pool @@ -29,6 +36,7 @@ services: redis: image: redis:8-alpine + logging: *log-limits ports: - "127.0.0.1:16379:6379/tcp" healthcheck: @@ -41,6 +49,7 @@ services: build: context: . dockerfile: Dockerfile.test + logging: *log-limits depends_on: timescaledb: condition: service_healthy diff --git a/docker-compose.yml b/docker-compose.yml index 59600ee..852005d 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,8 +1,15 @@ +x-log-limits: &log-limits + driver: json-file + options: + max-size: ${DOCKER_LOG_MAX_SIZE:-100m} + max-file: "${DOCKER_LOG_MAX_FILES:-5}" + services: timescaledb: image: timescale/timescaledb:latest-pg17 container_name: public-pool-timescaledb restart: unless-stopped + logging: *log-limits environment: POSTGRES_DB: ${DB_DATABASE:-public_pool} POSTGRES_USER: ${DB_USERNAME:-public_pool} @@ -33,6 +40,7 @@ services: image: redis:8-alpine container_name: public-pool-redis restart: unless-stopped + logging: *log-limits ports: - "127.0.0.1:${REDIS_PORT:-6379}:6379/tcp" volumes: @@ -49,6 +57,7 @@ services: context: . dockerfile: Dockerfile restart: unless-stopped + logging: *log-limits depends_on: timescaledb: condition: service_healthy @@ -83,6 +92,9 @@ services: PM2_ENABLED: ${PM2_ENABLED:-true} STRATUM_WORKERS: ${STRATUM_WORKERS:-2} STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330} + STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1} + STRATUM_SOCKET_TIMEOUT_MS: ${STRATUM_SOCKET_TIMEOUT_MS:-3600000} + STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS: ${STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS:-60000} STRATUM_SECURE: ${STRATUM_SECURE:-true} SECURE_STRATUM_PORTS: ${SECURE_STRATUM_PORTS:-4333,4332,4331,4330} STRATUM_MAX_CONNECTIONS_PER_LISTENER: ${STRATUM_MAX_CONNECTIONS_PER_LISTENER:-10000} diff --git a/ecosystem.config.js b/ecosystem.config.js index 5f82293..81d13ac 100644 --- a/ecosystem.config.js +++ b/ecosystem.config.js @@ -1,7 +1,14 @@ +const dockerLogConfig = { + out_file: '/dev/stdout', + error_file: '/dev/stderr', + merge_logs: true, +}; + module.exports = { apps: [ // API instance { + ...dockerLogConfig, name: 'api', script: './dist/main.js', instances: parseInt(process.env.API_WORKERS || '4', 10), @@ -16,6 +23,7 @@ module.exports = { }, // Master instance { + ...dockerLogConfig, name: 'master', script: './dist/main.js', instances: 1, @@ -29,6 +37,7 @@ module.exports = { }, // Worker instances { + ...dockerLogConfig, name: 'workers', script: './dist/main.js', instances: parseInt(process.env.STRATUM_WORKERS || '2', 10), diff --git a/full-setup/docker-compose-mainnet.yml b/full-setup/docker-compose-mainnet.yml index 2d438c0..440dc15 100644 --- a/full-setup/docker-compose-mainnet.yml +++ b/full-setup/docker-compose-mainnet.yml @@ -1,3 +1,9 @@ +x-log-limits: &log-limits + driver: json-file + options: + max-size: ${DOCKER_LOG_MAX_SIZE:-100m} + max-file: "${DOCKER_LOG_MAX_FILES:-5}" + services: bitcoin: container_name: bitcoin @@ -6,6 +12,7 @@ services: dockerfile: Dockerfile restart: unless-stopped stop_grace_period: 30s + logging: *log-limits networks: - bitcoin ports: @@ -22,6 +29,7 @@ services: image: timescale/timescaledb:latest-pg17 container_name: public-pool-mainnet-timescaledb restart: unless-stopped + logging: *log-limits networks: - bitcoin environment: @@ -52,6 +60,7 @@ services: image: redis:8-alpine container_name: public-pool-mainnet-redis restart: unless-stopped + logging: *log-limits networks: - bitcoin volumes: @@ -69,6 +78,7 @@ services: dockerfile: Dockerfile restart: unless-stopped stop_grace_period: 30s + logging: *log-limits depends_on: timescaledb: condition: service_healthy @@ -100,6 +110,9 @@ services: API_WORKERS: ${API_WORKERS:-4} PM2_ENABLED: "true" STRATUM_WORKERS: ${STRATUM_WORKERS:-2} + STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1} + STRATUM_SOCKET_TIMEOUT_MS: ${STRATUM_SOCKET_TIMEOUT_MS:-3600000} + STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS: ${STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS:-60000} networks: bitcoin: diff --git a/full-setup/docker-compose-regtest.yml b/full-setup/docker-compose-regtest.yml index 06c4e40..678cc1b 100644 --- a/full-setup/docker-compose-regtest.yml +++ b/full-setup/docker-compose-regtest.yml @@ -1,3 +1,9 @@ +x-log-limits: &log-limits + driver: json-file + options: + max-size: ${DOCKER_LOG_MAX_SIZE:-100m} + max-file: "${DOCKER_LOG_MAX_FILES:-5}" + services: bitcoin-regtest: container_name: bitcoin-regtest @@ -6,6 +12,7 @@ services: dockerfile: Dockerfile restart: unless-stopped stop_grace_period: 30s + logging: *log-limits networks: - bitcoin-regtest ports: @@ -24,6 +31,7 @@ services: image: timescale/timescaledb:latest-pg17 container_name: public-pool-regtest-timescaledb restart: unless-stopped + logging: *log-limits networks: - bitcoin-regtest environment: @@ -54,6 +62,7 @@ services: image: redis:8-alpine container_name: public-pool-regtest-redis restart: unless-stopped + logging: *log-limits networks: - bitcoin-regtest volumes: @@ -71,6 +80,7 @@ services: dockerfile: Dockerfile restart: unless-stopped stop_grace_period: 30s + logging: *log-limits depends_on: timescaledb-regtest: condition: service_healthy @@ -97,6 +107,9 @@ services: REDIS_URL: redis://redis-regtest:6379 PM2_ENABLED: "true" STRATUM_WORKERS: ${STRATUM_WORKERS:-2} + STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1} + STRATUM_SOCKET_TIMEOUT_MS: ${STRATUM_SOCKET_TIMEOUT_MS:-3600000} + STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS: ${STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS:-60000} networks: bitcoin-regtest: diff --git a/full-setup/docker-compose-testnet.yml b/full-setup/docker-compose-testnet.yml index 82819fc..6cf073e 100644 --- a/full-setup/docker-compose-testnet.yml +++ b/full-setup/docker-compose-testnet.yml @@ -1,3 +1,9 @@ +x-log-limits: &log-limits + driver: json-file + options: + max-size: ${DOCKER_LOG_MAX_SIZE:-100m} + max-file: "${DOCKER_LOG_MAX_FILES:-5}" + services: bitcoin-testnet: container_name: bitcoin-testnet @@ -6,6 +12,7 @@ services: dockerfile: Dockerfile restart: unless-stopped stop_grace_period: 30s + logging: *log-limits networks: - bitcoin-testnet ports: @@ -23,6 +30,7 @@ services: image: timescale/timescaledb:latest-pg17 container_name: public-pool-testnet-timescaledb restart: unless-stopped + logging: *log-limits networks: - bitcoin-testnet environment: @@ -53,6 +61,7 @@ services: image: redis:8-alpine container_name: public-pool-testnet-redis restart: unless-stopped + logging: *log-limits networks: - bitcoin-testnet volumes: @@ -70,6 +79,7 @@ services: dockerfile: Dockerfile restart: unless-stopped stop_grace_period: 30s + logging: *log-limits depends_on: timescaledb-testnet: condition: service_healthy @@ -96,6 +106,9 @@ services: REDIS_URL: redis://redis-testnet:6379 PM2_ENABLED: "true" STRATUM_WORKERS: ${STRATUM_WORKERS:-2} + STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1} + STRATUM_SOCKET_TIMEOUT_MS: ${STRATUM_SOCKET_TIMEOUT_MS:-3600000} + STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS: ${STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS:-60000} networks: bitcoin-testnet: diff --git a/src/ORM/_migrations/CurrentRoundWorkRollup1780899000000.ts b/src/ORM/_migrations/CurrentRoundWorkRollup1780899000000.ts new file mode 100644 index 0000000..9d60e78 --- /dev/null +++ b/src/ORM/_migrations/CurrentRoundWorkRollup1780899000000.ts @@ -0,0 +1,43 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +export class CurrentRoundWorkRollup1780899000000 implements MigrationInterface { + public name = 'CurrentRoundWorkRollup1780899000000'; + public transaction = false; + + public async up(queryRunner: QueryRunner): Promise { + await queryRunner.query(` + CREATE MATERIALIZED VIEW IF NOT EXISTS "accepted_share_block_10m" + WITH (timescaledb.continuous) AS + SELECT + time_bucket(INTERVAL '10 minutes', "acceptedAt") AS "bucket", + "blockHeight", + SUM("creditedDifficulty") AS "shares", + COUNT(*) AS "acceptedCount", + MAX("networkDifficulty") AS "networkDifficulty", + MAX("submissionDifficulty") AS "bestSubmissionDifficulty" + FROM "accepted_share_entity" + GROUP BY "bucket", "blockHeight" + WITH NO DATA + `); + + await queryRunner.query(` + SELECT add_continuous_aggregate_policy( + 'accepted_share_block_10m', + start_offset => INTERVAL '180 days', + end_offset => INTERVAL '1 minute', + schedule_interval => INTERVAL '1 minute', + if_not_exists => TRUE + ) + `); + + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS "IDX_accepted_share_block_10m_height_bucket" + ON "accepted_share_block_10m" ("blockHeight" DESC, "bucket" DESC) + `); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_block_10m_height_bucket"`); + await queryRunner.query(`DROP MATERIALIZED VIEW IF EXISTS "accepted_share_block_10m"`); + } +} diff --git a/src/ORM/_views/user-agent-report/user-agent-report.service.ts b/src/ORM/_views/user-agent-report/user-agent-report.service.ts index 85bfd58..ef67f24 100644 --- a/src/ORM/_views/user-agent-report/user-agent-report.service.ts +++ b/src/ORM/_views/user-agent-report/user-agent-report.service.ts @@ -5,7 +5,6 @@ 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 { @@ -23,27 +22,21 @@ export class UserAgentReportService { } public async getReport() { - const start = timingStart(); - const cachedReport = await timeAsync('user agent report shared cache read', () => this.redisMessagingService + const cachedReport = await this.redisMessagingService .getJsonCache(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') { - 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 await this.userAgentReport.find(); } - const rows = await this.refreshLiveReport(); - logTiming('user agent report getReport', start, { cache: 'miss', source: 'live-presence', rows: rows.length }); - return rows; + return await this.refreshLiveReport(); } public async refreshLiveReport() { @@ -60,8 +53,7 @@ export class UserAgentReportService { } private async buildLiveReport() { - const start = timingStart(); - const presences = await timeAsync('user agent report all presence load', () => this.redisMessagingService.getAllClientPresence()); + const presences = await this.redisMessagingService.getAllClientPresence(); const rows = new Map this.dataSource.query(` + const result = await 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'); } @@ -91,12 +90,7 @@ export class ClientStatisticsService { ORDER BY "label" `; - const result = await timeAsync('client statistics chart query', () => this.dataSource.query(query, params), { - filterSql, - params: params.length, - limit, - windowSql, - }); + const result = await this.dataSource.query(query, params); return result.map(res => { return { diff --git a/src/ORM/share-accounting/share-accounting.service.spec.ts b/src/ORM/share-accounting/share-accounting.service.spec.ts index f0c6d6a..4de1e78 100644 --- a/src/ORM/share-accounting/share-accounting.service.spec.ts +++ b/src/ORM/share-accounting/share-accounting.service.spec.ts @@ -153,6 +153,10 @@ describe('ShareAccountingService', () => { hashRateLastHour: 114532461.2, bestSubmissionDifficulty: 0, bestSubmissionDifficultyAt: null, + workSinceLastBlock: 0, + currentRoundAcceptedShares: 0, + currentRoundNetworkDifficulty: 0, + networkDifficultyPercent: 0, blockCandidateCount: 0, latestShareAt: '2026-06-07T12:10:00.000Z', protocolBreakdown: [], @@ -274,6 +278,11 @@ describe('ShareAccountingService', () => { .mockResolvedValueOnce([{ bestSubmissionDifficulty: '4096', bestSubmissionDifficultyAt: new Date('2026-06-07T12:19:00Z'), + }]) + .mockResolvedValueOnce([{ + currentRoundAcceptedShares: '11', + workSinceLastBlock: '352', + currentRoundNetworkDifficulty: '1000', }]), }; const service = new ShareAccountingService(repository as any, redis as any); @@ -284,12 +293,20 @@ describe('ShareAccountingService', () => { hashRateLast10Minutes: 1603451170.77, bestSubmissionDifficulty: 4096, bestSubmissionDifficultyAt: '2026-06-07T12:19:00.000Z', + workSinceLastBlock: 352, + currentRoundAcceptedShares: 11, + currentRoundNetworkDifficulty: 1000, + networkDifficultyPercent: 35.2, latestShareAt: '2026-06-07T12:20:00.000Z', })); expect(repository.query).toHaveBeenNthCalledWith( 3, expect.stringContaining('WHERE "blockHeight" > latest_found_block."height"'), ); + expect(repository.query).toHaveBeenNthCalledWith( + 4, + expect.stringContaining('FROM "accepted_share_block_10m"'), + ); }); it('should cache accounting summaries briefly to protect hot dashboard endpoints', async () => { diff --git a/src/ORM/share-accounting/share-accounting.service.ts b/src/ORM/share-accounting/share-accounting.service.ts index fc964f1..b4585ae 100644 --- a/src/ORM/share-accounting/share-accounting.service.ts +++ b/src/ORM/share-accounting/share-accounting.service.ts @@ -4,7 +4,6 @@ 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'; @@ -40,6 +39,10 @@ export interface ShareAccountingSummary { hashRateLastHour: number; bestSubmissionDifficulty: number; bestSubmissionDifficultyAt: string | null; + workSinceLastBlock: number; + currentRoundAcceptedShares: number; + currentRoundNetworkDifficulty: number; + networkDifficultyPercent: number; blockCandidateCount: number; latestShareAt: string | null; protocolBreakdown: { @@ -167,8 +170,8 @@ export class ShareAccountingService implements OnModuleDestroy { return summary; } - const [[liveWindow], [bestDifficultyRow]] = await Promise.all([ - timeAsync('share accounting live 10m pool query', () => this.acceptedShareRepository.query(` + const [[liveWindow], [bestDifficultyRow], [currentRoundRow]] = await Promise.all([ + this.acceptedShareRepository.query(` SELECT COUNT(*)::int AS "acceptedSharesLast10Minutes", COALESCE(SUM("creditedDifficulty"), 0)::float AS "creditedDifficultyLast10Minutes", @@ -176,8 +179,8 @@ export class ShareAccountingService implements OnModuleDestroy { 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(` + `), + this.acceptedShareRepository.query(` WITH latest_found_block AS ( SELECT COALESCE(MAX("height"), 0) AS "height" FROM "blocks_entity" @@ -189,8 +192,22 @@ export class ShareAccountingService implements OnModuleDestroy { WHERE "blockHeight" > latest_found_block."height" ORDER BY "submissionDifficulty" DESC, "acceptedAt" DESC LIMIT 1 - `)), + `), + this.acceptedShareRepository.query(` + WITH latest_found_block AS ( + SELECT COALESCE(MAX("height"), 0) AS "height" + FROM "blocks_entity" + ) + SELECT + COALESCE(SUM("acceptedCount"), 0)::int AS "currentRoundAcceptedShares", + COALESCE(SUM("shares"), 0)::float AS "workSinceLastBlock", + COALESCE(MAX("networkDifficulty"), 0)::float AS "currentRoundNetworkDifficulty" + FROM "accepted_share_block_10m", latest_found_block + WHERE "blockHeight" > latest_found_block."height" + `), ]); + const currentRoundNetworkDifficulty = this.toNumber(currentRoundRow?.currentRoundNetworkDifficulty); + const workSinceLastBlock = this.toNumber(currentRoundRow?.workSinceLastBlock); return { ...summary, @@ -201,6 +218,12 @@ export class ShareAccountingService implements OnModuleDestroy { bestSubmissionDifficultyAt: bestDifficultyRow?.bestSubmissionDifficultyAt == null ? null : new Date(bestDifficultyRow.bestSubmissionDifficultyAt).toISOString(), + workSinceLastBlock, + currentRoundAcceptedShares: this.toNumber(currentRoundRow?.currentRoundAcceptedShares), + currentRoundNetworkDifficulty, + networkDifficultyPercent: currentRoundNetworkDifficulty > 0 + ? this.roundPercent((workSinceLastBlock / currentRoundNetworkDifficulty) * 100) + : 0, latestShareAt: liveWindow?.latestShareAt == null ? summary.latestShareAt : new Date(liveWindow.latestShareAt).toISOString(), @@ -221,6 +244,10 @@ export class ShareAccountingService implements OnModuleDestroy { hashRateLastHour: 0, bestSubmissionDifficulty: 0, bestSubmissionDifficultyAt: null, + workSinceLastBlock: 0, + currentRoundAcceptedShares: 0, + currentRoundNetworkDifficulty: 0, + networkDifficultyPercent: 0, blockCandidateCount: 0, latestShareAt: null, protocolBreakdown: [], @@ -246,7 +273,7 @@ export class ShareAccountingService implements OnModuleDestroy { return summaries; } - const rows = await timeAsync('share accounting session summaries query', () => this.acceptedShareRepository.query(` + const rows = await this.acceptedShareRepository.query(` SELECT "clientId", MAX("bucket") AS "latestShareAt", @@ -254,7 +281,7 @@ export class ShareAccountingService implements OnModuleDestroy { FROM "accepted_share_10m" WHERE "clientId" = ANY($1::uuid[]) GROUP BY "clientId" - `, [uniqueClientIds]), { clientIds: uniqueClientIds.length }); + `, [uniqueClientIds]); rows.forEach(row => { summaries.set(row.clientId, { @@ -320,7 +347,7 @@ export class ShareAccountingService implements OnModuleDestroy { private async loadSummary(filter: AccountingFilter): Promise { const { whereSql, params } = this.buildWhereClause(filter); - const [summary] = await timeAsync('share accounting summary query', () => this.acceptedShareRepository.query(` + const [summary] = await this.acceptedShareRepository.query(` SELECT COALESCE(SUM("acceptedCount"), 0)::int AS "totalAcceptedShares", COALESCE(SUM("shares"), 0)::float AS "totalCreditedDifficulty", @@ -335,7 +362,7 @@ export class ShareAccountingService implements OnModuleDestroy { MAX("bucket") AS "latestShareAt" FROM "accepted_share_10m" ${whereSql} - `, params), { filter, params: params.length }); + `, params); return { totalAcceptedShares: this.toNumber(summary?.totalAcceptedShares), @@ -350,6 +377,10 @@ export class ShareAccountingService implements OnModuleDestroy { hashRateLastHour: this.toNumber(summary?.hashRateLastHour), bestSubmissionDifficulty: 0, bestSubmissionDifficultyAt: null, + workSinceLastBlock: 0, + currentRoundAcceptedShares: 0, + currentRoundNetworkDifficulty: 0, + networkDifficultyPercent: 0, blockCandidateCount: 0, latestShareAt: summary?.latestShareAt == null ? null @@ -436,6 +467,10 @@ export class ShareAccountingService implements OnModuleDestroy { return Number.isFinite(parsed) ? parsed : 0; } + private roundPercent(value: number): number { + return Math.round(value * 1_000_000) / 1_000_000; + } + private getSummaryCacheKey(filter: AccountingFilter): string { return JSON.stringify({ address: filter.address ?? null, diff --git a/src/app.controller.ts b/src/app.controller.ts index 73327ec..4d45040 100644 --- a/src/app.controller.ts +++ b/src/app.controller.ts @@ -13,7 +13,6 @@ 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 { @@ -36,28 +35,21 @@ export class AppController { @Get('info') public async info() { - const start = timingStart(); - - const CACHE_KEY = 'SITE_INFO'; 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; } const staleResult = await this.getCached(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; + return await this.refreshSiteInfo(null); } @@ -82,13 +74,13 @@ export class AppController { }; const [blockData, highScores, poolAuthority, userAgentReport] = await Promise.all([ - 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()), { + withInfoTimeout('found blocks', this.blocksService.getFoundBlocks(), staleInfo?.blockData ?? []), + withInfoTimeout('high scores', this.addressSettingsService.getHighScores(), staleInfo?.highScores ?? []), + withInfoTimeout('SV2 authority', this.stratumV2Service.getPoolAuthorityPublicKey(), { publicKey: staleInfo?.sv2?.poolAuthorityPublicKey ?? '', configured: staleInfo?.sv2?.authorityKeyConfigured ?? false }), - withInfoTimeout('user agent report', timeAsync('/api/info user agent report', () => this.userAgentReportService.getReport()), staleInfo?.userAgents ?? []), + withInfoTimeout('user agent report', this.userAgentReportService.getReport(), staleInfo?.userAgents ?? []), ]); const other: { @@ -139,42 +131,36 @@ export class AppController { @Get('info/accounting') public async infoAccounting() { - const start = timingStart(); const CACHE_KEY = 'SITE_ACCOUNTING'; 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 timeAsync('/api/info/accounting getPoolSummary', () => this.shareAccountingService.getPoolSummary()); + const data = await 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, 15 * 1000); if (cachedResult != null) { - logTiming('GET /api/pool', start, { cache: 'hit' }); return cachedResult; } - const userAgents = await timeAsync('/api/pool user agent report', () => this.userAgentReportService.getReport()); + const userAgents = await this.userAgentReportService.getReport(); const totalHashRate = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.totalHashRate), 0); const totalMiners = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.count), 0); const blockHeight = this.bitcoinRpcService.miningInfo.blocks; - const blocksFound = await timeAsync('/api/pool found blocks', () => this.blocksService.getFoundBlocks()); + const blocksFound = await this.blocksService.getFoundBlocks(); const data = { totalHashRate, @@ -187,36 +173,28 @@ 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, 10 * 60 * 1000); if (cachedResult != null) { - logTiming('GET /api/info/chart', start, { cache: 'hit' }); return cachedResult; } - const chartData = await timeAsync('/api/info/chart query', () => this.clientStatisticsService.getChartDataForSite()); + const chartData = await 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; diff --git a/src/controllers/client/client.controller.ts b/src/controllers/client/client.controller.ts index ef17041..31f1d6f 100644 --- a/src/controllers/client/client.controller.ts +++ b/src/controllers/client/client.controller.ts @@ -5,7 +5,6 @@ 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') @@ -22,15 +21,13 @@ export class ClientController { @Get(':address') async getClientInfo(@Param('address') address: string) { - const start = timingStart(); - - 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 workers = await this.redisMessagingService.getClientPresenceByAddress(address); + const sessionSummaries = await this.shareAccountingService.getSessionSummaries(workers.map(worker => worker.clientId)); 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 }); + : await this.addressSettingsService.getSettings(address, false); + const accounting = await this.shareAccountingService.getAddressSummary(address); const bestDifficulty = addressSettings?.bestDifficulty ?? workers.reduce((best, worker) => { return Math.max(best, Number(worker.bestDifficulty ?? 0)); }, 0); @@ -57,23 +54,17 @@ 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 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; + return await this.clientStatisticsService.getChartDataForAddress(address); } @Get(':address/:workerName') async getWorkerGroupInfo(@Param('address') address: string, @Param('workerName') workerName: string) { - const start = timingStart(); - - const addressWorkers = await timeAsync('/api/client/:address/:workerName address presence', () => this.redisMessagingService.getClientPresenceByAddress(address), { address, workerName }); + const addressWorkers = await this.redisMessagingService.getClientPresenceByAddress(address); const workers = addressWorkers .filter(worker => worker.clientName === workerName); @@ -84,8 +75,8 @@ export class ClientController { return pre; }, 0); - 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 chartData = await this.clientStatisticsService.getChartDataForGroup(address, workerName); + const accounting = await this.shareAccountingService.getWorkerGroupSummary(address, workerName); const response = { name: workerName, @@ -94,19 +85,16 @@ export class ClientController { 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 addressWorkers = await timeAsync('/api/client/:address/:workerName/:sessionId address presence', () => this.redisMessagingService.getClientPresenceByAddress(address), { address, workerName, sessionId }); + const addressWorkers = await this.redisMessagingService.getClientPresenceByAddress(address); const presenceWorker = addressWorkers .find(worker => worker.clientName === workerName && worker.sessionId === sessionId); const worker = presenceWorker == null - ? await timeAsync('/api/client/:address/:workerName/:sessionId DB fallback', () => this.clientService.getBySessionId(address, workerName, sessionId), { address, workerName, sessionId }) + ? await this.clientService.getBySessionId(address, workerName, sessionId) : { id: presenceWorker.clientId, sessionId: presenceWorker.sessionId, @@ -115,11 +103,10 @@ 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 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 }); + const chartData = await this.clientStatisticsService.getChartDataForSession(worker.id); + const accounting = await this.shareAccountingService.getSessionSummary(worker.id); const response = { sessionId: worker.sessionId, @@ -129,7 +116,6 @@ export class ClientController { 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; } } diff --git a/src/database.config.ts b/src/database.config.ts index 92b0771..871139e 100644 --- a/src/database.config.ts +++ b/src/database.config.ts @@ -8,6 +8,7 @@ import { AcceptedShareIndex1780862400000 } from './ORM/_migrations/AcceptedShare import { AcceptedShareRollupIndexes1780865400000 } from './ORM/_migrations/AcceptedShareRollupIndexes1780865400000'; import { PoolAccountingDashboardIndexes1780867200000 } from './ORM/_migrations/PoolAccountingDashboardIndexes1780867200000'; import { CurrentRoundBestShareIndex1780897600000 } from './ORM/_migrations/CurrentRoundBestShareIndex1780897600000'; +import { CurrentRoundWorkRollup1780899000000 } from './ORM/_migrations/CurrentRoundWorkRollup1780899000000'; 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'; @@ -34,6 +35,7 @@ export const databaseMigrations = [ AcceptedShareRollupIndexes1780865400000, PoolAccountingDashboardIndexes1780867200000, CurrentRoundBestShareIndex1780897600000, + CurrentRoundWorkRollup1780899000000, ]; export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions { diff --git a/src/models/StratumV1Client.spec.ts b/src/models/StratumV1Client.spec.ts index 6b8c86b..9134888 100644 --- a/src/models/StratumV1Client.spec.ts +++ b/src/models/StratumV1Client.spec.ts @@ -186,6 +186,17 @@ describe('StratumV1Client', () => { }); + it('should disable application idle timeout after Stratum initialization', async () => { + const setTimeoutSpy = jest.spyOn(socket, 'setTimeout').mockImplementation(() => socket); + jest.spyOn(client as any, 'write').mockImplementation(() => Promise.resolve(true)); + + emitMessage(MockRecording1.MINING_SUBSCRIBE); + emitMessage(MockRecording1.MINING_AUTHORIZE); + await new Promise((r) => setTimeout(r, 100)); + + expect(setTimeoutSpy).toHaveBeenCalledWith(0); + }); + it('should block non-compliant user agents on subscribe without allocating a session', async () => { (configService.get as jest.Mock).mockImplementation((key: string) => { switch (key) { @@ -282,6 +293,26 @@ describe('StratumV1Client', () => { expect(socket.write).toHaveBeenCalledWith(`{"id":null,"method":"mining.set_difficulty","params":[512]}\n`, expect.any(Function)); }); + it('should clamp suggested difficulty to the configured minimum', async () => { + (configService.get as jest.Mock).mockImplementation((key: string) => { + switch (key) { + case 'STRATUM_MIN_DIFFICULTY': + return '1'; + case 'DEV_FEE_ADDRESS': + return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4'; + case 'NETWORK': + return 'testnet'; + } + return null; + }); + jest.spyOn(socket, 'write').mockImplementation((data) => true); + + emitMessage(`{"id":4,"method":"mining.suggest_difficulty","params":[0]}`); + await new Promise((r) => setTimeout(r, 1)); + + expect(socket.write).toHaveBeenCalledWith(`{"id":null,"method":"mining.set_difficulty","params":[1]}\n`, expect.any(Function)); + }); + it('should set difficulty', async () => { jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true)); diff --git a/src/models/StratumV1Client.ts b/src/models/StratumV1Client.ts index 2da937d..7a867ef 100644 --- a/src/models/StratumV1Client.ts +++ b/src/models/StratumV1Client.ts @@ -32,6 +32,7 @@ import { StratumV1ClientStatistics } from './StratumV1ClientStatistics'; const TRUE_DIFF_ONE = 2.695953529101131e67; const BLOCKED_USER_AGENT_LOG_INTERVAL_MS = 60 * 1000; const VALIDATION_ERROR_LOG_INTERVAL_MS = 60 * 1000; +const DEFAULT_MIN_DIFFICULTY = 0.001; export class StratumV1Client { private static blockedUserAgentLogState = new Map(); @@ -159,7 +160,7 @@ export class StratumV1Client { if (this.sessionStart == null) { this.sessionStart = new Date(); - this.statistics = new StratumV1ClientStatistics(); + this.statistics = new StratumV1ClientStatistics(this.getMinimumDifficulty()); this.extraNonceAndSessionId = this.getRandomHexString(); //console.log(`New client ID: : ${this.extraNonceAndSessionId}, ${this.socket.remoteAddress}:${this.socket.remotePort}`); } @@ -281,7 +282,7 @@ export class StratumV1Client { if (errors.length === 0) { this.clientSuggestedDifficulty = suggestDifficultyMessage; - this.sessionDifficulty = suggestDifficultyMessage.suggestedDifficulty; + this.sessionDifficulty = this.clampDifficulty(suggestDifficultyMessage.suggestedDifficulty); const success = await this.write(JSON.stringify(this.clientSuggestedDifficulty.response(this.sessionDifficulty)) + '\n'); if (!success) { return; @@ -369,6 +370,7 @@ export class StratumV1Client { private async initStratum() { this.stratumInitialized = true; + this.socket.setTimeout(0); if (this.isBlockedUserAgent(this.clientSubscription.userAgent)) { this.logBlockedUserAgent(this.clientSubscription.userAgent); @@ -376,12 +378,6 @@ export class StratumV1Client { return; } - switch (this.clientSubscription.userAgent) { - case 'cpuminer': { - this.sessionDifficulty = 0.1; - } - } - if (this.clientSuggestedDifficulty == null) { //console.log(`Setting difficulty to ${this.sessionDifficulty}`) const setDifficulty = JSON.stringify(new SuggestDifficulty().response(this.sessionDifficulty)); @@ -688,7 +684,7 @@ export class StratumV1Client { } private async checkDifficulty() { - const targetDiff = this.statistics.getSuggestedDifficulty(this.sessionDifficulty); + const targetDiff = this.clampDifficulty(this.statistics.getSuggestedDifficulty(this.sessionDifficulty)); if (targetDiff == null) { return; } @@ -824,6 +820,27 @@ export class StratumV1Client { return ` sample=${values.join(',')}`; } + private clampDifficulty(difficulty: number | null): number | null { + if (difficulty == null || !Number.isFinite(difficulty)) { + return null; + } + const configuredMinimum = this.getConfiguredMinimumDifficulty(); + return configuredMinimum == null ? difficulty : Math.max(difficulty, configuredMinimum); + } + + private getMinimumDifficulty(): number { + return this.getConfiguredMinimumDifficulty() ?? DEFAULT_MIN_DIFFICULTY; + } + + private getConfiguredMinimumDifficulty(): number | null { + const configured = parseFloat( + this.configService.get('STRATUM_MIN_DIFFICULTY') + ?? process.env.STRATUM_MIN_DIFFICULTY + ?? '', + ); + return Number.isFinite(configured) && configured > 0 ? configured : null; + } + private closeSocket() { this.connectionClosed = true; if (!this.socket.destroyed) { diff --git a/src/models/StratumV1ClientStatistics.spec.ts b/src/models/StratumV1ClientStatistics.spec.ts index 6668594..0b837f7 100644 --- a/src/models/StratumV1ClientStatistics.spec.ts +++ b/src/models/StratumV1ClientStatistics.spec.ts @@ -62,4 +62,11 @@ describe('StratumV1ClientStatistics', () => { expect(statistics.getSuggestedDifficulty(128)).toBe(16); }); + + it('should not suggest a difficulty below the configured minimum', () => { + statistics = new StratumV1ClientStatistics(1); + jest.setSystemTime(new Date('2026-05-06T12:06:00Z')); + + expect(statistics.getSuggestedDifficulty(1)).toBe(1); + }); }); diff --git a/src/models/StratumV1ClientStatistics.ts b/src/models/StratumV1ClientStatistics.ts index 86cee6b..e56c363 100644 --- a/src/models/StratumV1ClientStatistics.ts +++ b/src/models/StratumV1ClientStatistics.ts @@ -1,7 +1,7 @@ import { ClientEntity } from '../ORM/client/client.entity'; const CACHE_SIZE = 30; -const MIN_DIFF = 0.001; +const DEFAULT_MIN_DIFF = 0.001; export class StratumV1ClientStatistics { public targetSubmitShareEveryNSeconds: number = 30; @@ -11,7 +11,7 @@ export class StratumV1ClientStatistics { private submissionCache: { time: Date, difficulty: number }[] = []; private submissionCacheDifficultySum = 0; - constructor() { + constructor(private readonly minDifficulty = DEFAULT_MIN_DIFF) { this.submissionCacheStart = new Date(); } @@ -67,8 +67,8 @@ export class StratumV1ClientStatistics { if (val === 0) { return null; } - if (val < MIN_DIFF) { - return MIN_DIFF; + if (val < this.minDifficulty) { + return this.minDifficulty; } let x = val | (val >> 1); x = x | (x >> 2); @@ -77,8 +77,8 @@ export class StratumV1ClientStatistics { x = x | (x >> 16); x = x | (x >> 32); const res = x - (x >> 1); - if (res == 0 && val * 100 < MIN_DIFF) { - return MIN_DIFF; + if (res == 0 && val * 100 < this.minDifficulty) { + return this.minDifficulty; } if (res == 0) { return this.nearestPowerOfTwo(val * 100) / 100; diff --git a/src/models/StratumV2Client.ts b/src/models/StratumV2Client.ts index 8c86f08..de1a907 100644 --- a/src/models/StratumV2Client.ts +++ b/src/models/StratumV2Client.ts @@ -58,6 +58,7 @@ import { import { Sv2NoiseSession } from './sv2/sv2-noise'; const DEFAULT_START_DIFFICULTY = 100000; +const DEFAULT_MIN_DIFFICULTY = 0.001; const DEFAULT_TARGET_SHARES_PER_MINUTE = 2; const DEFAULT_DIFFICULTY_CHECK_INTERVAL_MS = 60 * 1000; const FIXED_STANDARD_EXTRANONCE2 = '0000000000000000'; @@ -141,7 +142,7 @@ export class StratumV2Client { this.sessionDifficulty = this.getInitialDifficulty(); this.targetSharesPerMinute = this.getTargetSharesPerMinute(); this.difficultyCheckIntervalMs = this.getDifficultyCheckIntervalMs(); - this.statistics = new StratumV1ClientStatistics(); + this.statistics = new StratumV1ClientStatistics(this.getMinimumDifficulty()); this.statistics.targetSubmitShareEveryNSeconds = 60 / this.targetSharesPerMinute; this.network = this.getNetwork(); @@ -357,10 +358,10 @@ export class StratumV2Client { this.targetSharesPerMinute, ); if (Number.isFinite(calculatedDifficulty) && calculatedDifficulty > 0) { - channelDifficulty = calculatedDifficulty; + channelDifficulty = this.clampDifficulty(calculatedDifficulty); } } - channelDifficulty = DifficultyUtils.clampDifficultyToMaxTarget(channelDifficulty, message.maxTarget); + channelDifficulty = this.clampDifficulty(DifficultyUtils.clampDifficultyToMaxTarget(channelDifficulty, message.maxTarget)); this.sessionDifficulty = channelDifficulty; const channel: ChannelState = { @@ -385,6 +386,7 @@ export class StratumV2Client { } await this.ensureClientEntity(); + this.disableApplicationIdleTimeout(); await this.sendFrame( Sv2MsgType.OPEN_STANDARD_MINING_CHANNEL_SUCCESS, serializeOpenStandardMiningChannelSuccess({ @@ -449,10 +451,10 @@ export class StratumV2Client { this.targetSharesPerMinute, ); if (Number.isFinite(calculatedDifficulty) && calculatedDifficulty > 0) { - channelDifficulty = calculatedDifficulty; + channelDifficulty = this.clampDifficulty(calculatedDifficulty); } } - channelDifficulty = DifficultyUtils.clampDifficultyToMaxTarget(channelDifficulty, message.maxTarget); + channelDifficulty = this.clampDifficulty(DifficultyUtils.clampDifficultyToMaxTarget(channelDifficulty, message.maxTarget)); this.sessionDifficulty = channelDifficulty; const channel: ChannelState = { @@ -477,6 +479,7 @@ export class StratumV2Client { } await this.ensureClientEntity(); + this.disableApplicationIdleTimeout(); await this.sendFrame( Sv2MsgType.OPEN_EXTENDED_MINING_CHANNEL_SUCCESS, serializeOpenExtendedMiningChannelSuccess({ @@ -797,10 +800,10 @@ export class StratumV2Client { this.targetSharesPerMinute, ); if (Number.isFinite(nextDifficulty) && nextDifficulty > 0) { - channel.sessionDifficulty = DifficultyUtils.clampDifficultyToMaxTarget( + channel.sessionDifficulty = this.clampDifficulty(DifficultyUtils.clampDifficultyToMaxTarget( nextDifficulty, channel.declaredMaxTarget, - ); + )); await this.sendSetTarget(channel); } } @@ -876,17 +879,17 @@ export class StratumV2Client { } private async checkDifficulty(): Promise { - const targetDiff = this.statistics.getSuggestedDifficulty(this.sessionDifficulty); + const targetDiff = this.clampDifficulty(this.statistics.getSuggestedDifficulty(this.sessionDifficulty)); if (targetDiff == null || targetDiff === this.sessionDifficulty || !Number.isFinite(targetDiff)) { return; } this.sessionDifficulty = targetDiff; for (const channel of this.channels.values()) { - channel.sessionDifficulty = DifficultyUtils.clampDifficultyToMaxTarget( + channel.sessionDifficulty = this.clampDifficulty(DifficultyUtils.clampDifficultyToMaxTarget( targetDiff, channel.declaredMaxTarget, - ); + )); await this.sendSetTarget(channel); } @@ -1191,7 +1194,35 @@ export class StratumV2Client { ?? this.configService.get('STRATUM_START_DIFFICULTY') ?? '', ); - return Number.isFinite(configured) && configured > 0 ? configured : DEFAULT_START_DIFFICULTY; + const difficulty = Number.isFinite(configured) && configured > 0 ? configured : DEFAULT_START_DIFFICULTY; + return this.clampDifficulty(difficulty); + } + + private getMinimumDifficulty(): number { + return this.getConfiguredMinimumDifficulty() ?? DEFAULT_MIN_DIFFICULTY; + } + + private getConfiguredMinimumDifficulty(): number | null { + const configured = parseFloat( + this.configService.get('STRATUM_MIN_DIFFICULTY') + ?? process.env.STRATUM_MIN_DIFFICULTY + ?? '', + ); + return Number.isFinite(configured) && configured > 0 ? configured : null; + } + + private clampDifficulty(difficulty: number | null): number | null { + if (difficulty == null || !Number.isFinite(difficulty)) { + return null; + } + const configuredMinimum = this.getConfiguredMinimumDifficulty(); + return configuredMinimum == null ? difficulty : Math.max(difficulty, configuredMinimum); + } + + private disableApplicationIdleTimeout(): void { + if (typeof this.socket.setTimeout === 'function') { + this.socket.setTimeout(0); + } } private getTargetSharesPerMinute(): number { diff --git a/src/services/redis-messaging.service.ts b/src/services/redis-messaging.service.ts index bc222b4..123e55b 100644 --- a/src/services/redis-messaging.service.ts +++ b/src/services/redis-messaging.service.ts @@ -4,7 +4,6 @@ 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'; @@ -264,10 +263,8 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy { } private async getPresenceFromSet(setKey: string): Promise { - 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 []; } @@ -297,12 +294,6 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy { } } - logTiming('redis presence set load', start, { - setKey, - clientIds: clientIds.length, - presences: presences.length, - staleClientIds: staleClientIds.length, - }); return presences; } diff --git a/src/services/stratum-v1.service.spec.ts b/src/services/stratum-v1.service.spec.ts index 057d0e6..c769e9f 100644 --- a/src/services/stratum-v1.service.spec.ts +++ b/src/services/stratum-v1.service.spec.ts @@ -8,6 +8,8 @@ describe('StratumV1Service', () => { const originalBackpressureEnabled = process.env.STRATUM_BACKPRESSURE_ENABLED; const originalMaxConnectionsPerListener = process.env.STRATUM_MAX_CONNECTIONS_PER_LISTENER; const originalTlsHandshakeTimeoutMs = process.env.STRATUM_TLS_HANDSHAKE_TIMEOUT_MS; + const originalSocketTimeoutMs = process.env.STRATUM_SOCKET_TIMEOUT_MS; + const originalTcpKeepAliveInitialDelayMs = process.env.STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS; let service: StratumV1Service; let clientService; @@ -57,6 +59,8 @@ describe('StratumV1Service', () => { restoreEnv('STRATUM_BACKPRESSURE_ENABLED', originalBackpressureEnabled); restoreEnv('STRATUM_MAX_CONNECTIONS_PER_LISTENER', originalMaxConnectionsPerListener); restoreEnv('STRATUM_TLS_HANDSHAKE_TIMEOUT_MS', originalTlsHandshakeTimeoutMs); + restoreEnv('STRATUM_SOCKET_TIMEOUT_MS', originalSocketTimeoutMs); + restoreEnv('STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS', originalTcpKeepAliveInitialDelayMs); consoleLogSpy.mockRestore(); consoleWarnSpy.mockRestore(); jest.useRealTimers(); @@ -164,6 +168,30 @@ describe('StratumV1Service', () => { expect((service as any).getTlsHandshakeTimeoutMs()).toBe(5000); }); + it('should keep quiet miners connected for one hour by default', () => { + delete process.env.STRATUM_SOCKET_TIMEOUT_MS; + + expect((service as any).getSocketTimeoutMs()).toBe(1000 * 60 * 60); + }); + + it('should allow configuring Stratum socket idle timeout', () => { + process.env.STRATUM_SOCKET_TIMEOUT_MS = '7200000'; + + expect((service as any).getSocketTimeoutMs()).toBe(7200000); + }); + + it('should enable TCP keepalive quickly by default', () => { + delete process.env.STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS; + + expect((service as any).getTcpKeepAliveInitialDelayMs()).toBe(60000); + }); + + it('should allow configuring TCP keepalive initial delay', () => { + process.env.STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS = '30000'; + + expect((service as any).getTcpKeepAliveInitialDelayMs()).toBe(30000); + }); + it('should detect JSON-RPC as Stratum V1', () => { const firstChunk = Buffer.from('{"id":1,"method":"mining.subscribe","params":[]}\n'); diff --git a/src/services/stratum-v1.service.ts b/src/services/stratum-v1.service.ts index 1465b83..8d53673 100644 --- a/src/services/stratum-v1.service.ts +++ b/src/services/stratum-v1.service.ts @@ -35,6 +35,8 @@ const DEFAULT_BACKPRESSURE_RESUME_RSS_MB = 2000; const DEFAULT_BACKPRESSURE_HEALTHY_CHECKS = 3; const DEFAULT_MAX_CONNECTIONS_PER_LISTENER = 10000; const DEFAULT_TLS_HANDSHAKE_TIMEOUT_MS = 10000; +const DEFAULT_SOCKET_TIMEOUT_MS = 1000 * 60 * 60; +const DEFAULT_TCP_KEEPALIVE_INITIAL_DELAY_MS = 1000 * 60; @@ -118,8 +120,8 @@ export class StratumV1Service implements OnModuleInit { private createSocketServer(): Server { const server = new Server(async (socket: Socket) => { - // Set 15-minute timeout - socket.setTimeout(1000 * 60 * 15); + socket.setTimeout(this.getSocketTimeoutMs()); + socket.setKeepAlive(true, this.getTcpKeepAliveInitialDelayMs()); let client: StratumV1Client | StratumV2Client = null; let protocol: 'v1' | 'v2' | null = null; @@ -235,8 +237,8 @@ export class StratumV1Service implements OnModuleInit { }; const server = createServer(tlsOptions, async (socket: TLSSocket) => { - // Set 15-minute timeout - socket.setTimeout(1000 * 60 * 15); + socket.setTimeout(this.getSocketTimeoutMs()); + socket.setKeepAlive(true, this.getTcpKeepAliveInitialDelayMs()); const client = this.createV1Client(socket); @@ -421,6 +423,14 @@ export class StratumV1Service implements OnModuleInit { return this.getPositiveIntegerEnv('STRATUM_TLS_HANDSHAKE_TIMEOUT_MS', DEFAULT_TLS_HANDSHAKE_TIMEOUT_MS); } + private getSocketTimeoutMs() { + return this.getPositiveIntegerEnv('STRATUM_SOCKET_TIMEOUT_MS', DEFAULT_SOCKET_TIMEOUT_MS); + } + + private getTcpKeepAliveInitialDelayMs() { + return this.getPositiveIntegerEnv('STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS', DEFAULT_TCP_KEEPALIVE_INITIAL_DELAY_MS); + } + private detectProtocol(firstChunk: Buffer): 'v1' | 'v2' | null { if (firstChunk.length === 0) { return null; diff --git a/src/services/stratum-v2.service.ts b/src/services/stratum-v2.service.ts index e240d89..286c550 100644 --- a/src/services/stratum-v2.service.ts +++ b/src/services/stratum-v2.service.ts @@ -22,6 +22,9 @@ import { NotificationService } from './notification.service'; import { RedisMessagingService } from './redis-messaging.service'; import { StratumV1JobsService } from './stratum-v1-jobs.service'; +const DEFAULT_SOCKET_TIMEOUT_MS = 1000 * 60 * 60; +const DEFAULT_TCP_KEEPALIVE_INITIAL_DELAY_MS = 1000 * 60; + @Injectable() export class StratumV2Service implements OnModuleInit { private readonly servers: Server[] = []; @@ -171,6 +174,7 @@ export class StratumV2Service implements OnModuleInit { private startSocketServer(port: number): void { const server = new Server((socket: Socket) => { socket.setTimeout(this.getSocketTimeoutMs()); + socket.setKeepAlive(true, this.getTcpKeepAliveInitialDelayMs()); socket.setNoDelay(true); let client: StratumV2Client = null; @@ -213,7 +217,21 @@ export class StratumV2Service implements OnModuleInit { } private getSocketTimeoutMs(): number { - const configured = parseInt(this.configService.get('STRATUM_V2_SOCKET_TIMEOUT_MS') ?? '', 10); - return Number.isFinite(configured) && configured > 0 ? configured : 1000 * 60 * 15; + const configured = parseInt( + this.configService.get('STRATUM_V2_SOCKET_TIMEOUT_MS') + ?? this.configService.get('STRATUM_SOCKET_TIMEOUT_MS') + ?? '', + 10, + ); + return Number.isFinite(configured) && configured > 0 ? configured : DEFAULT_SOCKET_TIMEOUT_MS; + } + + private getTcpKeepAliveInitialDelayMs(): number { + const configured = parseInt( + this.configService.get('STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS') + ?? '', + 10, + ); + return Number.isFinite(configured) && configured > 0 ? configured : DEFAULT_TCP_KEEPALIVE_INITIAL_DELAY_MS; } } diff --git a/src/utils/timing.utils.ts b/src/utils/timing.utils.ts deleted file mode 100644 index d817f41..0000000 --- a/src/utils/timing.utils.ts +++ /dev/null @@ -1,51 +0,0 @@ -import { performance } from 'perf_hooks'; - -type TimingMetadata = Record | (() => Record); - -const DEFAULT_TIMING_LOG_MS = 250; - -export function timingStart(): number { - return performance.now(); -} - -export async function timeAsync( - label: string, - work: () => Promise, - metadata?: TimingMetadata, -): Promise { - 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 | null { - if (metadata == null) { - return null; - } - - try { - return typeof metadata === 'function' ? metadata() : metadata; - } catch (error) { - return { metadataError: error instanceof Error ? error.message : String(error) }; - } -} diff --git a/test/timescale-redis.integration-spec.ts b/test/timescale-redis.integration-spec.ts index 4cf70a0..3ae3b86 100644 --- a/test/timescale-redis.integration-spec.ts +++ b/test/timescale-redis.integration-spec.ts @@ -64,13 +64,14 @@ describe('TimescaleDB and Redis integration', () => { const aggregates = await dataSource.query(` SELECT view_name FROM timescaledb_information.continuous_aggregates - WHERE view_name IN ('accepted_share_10m', 'accepted_share_1h', 'accepted_share_1d') + WHERE view_name IN ('accepted_share_10m', 'accepted_share_1h', 'accepted_share_1d', 'accepted_share_block_10m') ORDER BY view_name `); expect(aggregates.map(row => row.view_name)).toEqual([ 'accepted_share_10m', 'accepted_share_1d', 'accepted_share_1h', + 'accepted_share_block_10m', ]); const legacyTables = await dataSource.query(` @@ -166,6 +167,7 @@ describe('TimescaleDB and Redis integration', () => { expect(Number(rows[0].last)).toBeGreaterThan(Number(rows[0].first)); await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_10m', NULL, NULL)`); + await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_block_10m', NULL, NULL)`); const aggregateRows = await dataSource.query(` SELECT "shares"::float AS shares, "acceptedCount"::int AS "acceptedCount" FROM accepted_share_10m @@ -175,6 +177,19 @@ describe('TimescaleDB and Redis integration', () => { expect(aggregateRows).toEqual(expect.arrayContaining([ expect.objectContaining({ shares: 96, acceptedCount: 2 }), ])); + + const blockAggregateRows = await dataSource.query(` + SELECT + "shares"::float AS shares, + "acceptedCount"::int AS "acceptedCount", + "networkDifficulty"::float AS "networkDifficulty" + FROM accepted_share_block_10m + WHERE "blockHeight" = $1 + `, [900000]); + + expect(blockAggregateRows).toEqual(expect.arrayContaining([ + expect.objectContaining({ shares: 96, acceptedCount: 2, networkDifficulty: 100000 }), + ])); }); it('should exclude soft-deleted clients from the user-agent report', async () => {