diff --git a/.env.example b/.env.example index 4c3aae0..431003c 100644 --- a/.env.example +++ b/.env.example @@ -100,6 +100,7 @@ SHARE_ACCOUNTING_MAX_QUEUE_SIZE=50000 SHARE_ACCOUNTING_SUMMARY_CACHE_MS=2500 SHARE_ACCOUNTING_SUMMARY_CACHE_MAX=10000 CLIENT_PRESENCE_ENABLED=true +CLIENT_HASHRATE_PERSIST_INTERVAL_MS=60000 SHARE_ROLLUP_ENABLED=true SHARE_ROLLUP_INTERVAL_MS=60000 SHARE_ROLLUP_SAFETY_LAG_SECONDS=30 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 5f869b5..ce42c88 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 @@ -62,6 +62,9 @@ export class UserAgentReportService { private async buildLiveReport() { const presences = await this.redisMessagingService.getAllClientPresence(); const activePresences = await this.filterActivePresences(presences); + if (activePresences.length === 0) { + return await this.userAgentReport.find(); + } const rows = new Map { return await this.clientRepository.count(); } diff --git a/src/models/StratumV1Client.ts b/src/models/StratumV1Client.ts index 4065f75..a5001e1 100644 --- a/src/models/StratumV1Client.ts +++ b/src/models/StratumV1Client.ts @@ -37,6 +37,7 @@ 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; +const DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS = 60 * 1000; export class StratumV1Client { private static blockedUserAgentLogState = new Map(); @@ -68,6 +69,7 @@ export class StratumV1Client { private buffer: string = ''; private connectionClosed = false; private lastSentMiningJobTimestamp: number = null; + private lastHashRatePersistedAt = 0; private miningSubmissionHashes = new Set() @@ -704,6 +706,7 @@ export class StratumV1Client { const now = new Date(); this.clientEntity.updatedAt = now; this.clientEntity.hashRate = this.statistics.hashRate; + await this.persistClientHashRate(now); await this.updateClientPresence(now); } catch (e) { @@ -852,6 +855,39 @@ export class StratumV1Client { }); } + private async persistClientHashRate(now: Date): Promise { + if (this.clientEntity?.id == null) { + return; + } + + const hashRate = Number(this.statistics?.hashRate ?? 0); + if (!Number.isFinite(hashRate) || hashRate <= 0) { + return; + } + + const intervalMs = this.getHashRatePersistIntervalMs(); + const nowMs = now.getTime(); + if (this.lastHashRatePersistedAt > 0 && nowMs - this.lastHashRatePersistedAt < intervalMs) { + return; + } + + this.lastHashRatePersistedAt = nowMs; + try { + await this.clientService.updateHashRate(this.clientEntity.id, hashRate, now); + } catch (error) { + console.error(`Failed to persist SV1 client hashrate: ${error.message}`); + } + } + + private getHashRatePersistIntervalMs(): number { + const configured = Number(this.configService.get('CLIENT_HASHRATE_PERSIST_INTERVAL_MS')); + if (Number.isFinite(configured) && configured >= 0) { + return configured; + } + + return DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS; + } + private getValidationErrorSignature(errors: ValidationError[]): string { if (errors.length === 0) { return 'unknown'; diff --git a/src/models/StratumV2Client.ts b/src/models/StratumV2Client.ts index 8e820de..1b1b9ae 100644 --- a/src/models/StratumV2Client.ts +++ b/src/models/StratumV2Client.ts @@ -70,6 +70,7 @@ 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 DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS = 60 * 1000; const FIXED_STANDARD_EXTRANONCE2 = '0000000000000000'; const RETIRED_EXTENDED_JOB_RETENTION_MS = 5 * 60 * 1000; const SV2_AUTH_FAILURE_LOG_INTERVAL_MS = 60 * 1000; @@ -133,6 +134,7 @@ export class StratumV2Client { private creatingEntity: Promise = null; private readonly firstChunkSummary: string; private workSelectionEnabled = false; + private lastHashRatePersistedAt = 0; constructor( private readonly socket: Socket, @@ -934,6 +936,7 @@ export class StratumV2Client { const now = new Date(); this.clientEntity.updatedAt = now; this.clientEntity.hashRate = this.statistics.hashRate; + await this.persistClientHashRate(now); await this.updateClientPresence(now); if (submissionDifficulty > this.clientEntity.bestDifficulty) { @@ -1343,6 +1346,39 @@ export class StratumV2Client { } } + private async persistClientHashRate(now: Date): Promise { + if (this.clientEntity?.id == null) { + return; + } + + const hashRate = Number(this.statistics?.hashRate ?? 0); + if (!Number.isFinite(hashRate) || hashRate <= 0) { + return; + } + + const intervalMs = this.getHashRatePersistIntervalMs(); + const nowMs = now.getTime(); + if (this.lastHashRatePersistedAt > 0 && nowMs - this.lastHashRatePersistedAt < intervalMs) { + return; + } + + this.lastHashRatePersistedAt = nowMs; + try { + await this.clientService.updateHashRate(this.clientEntity.id, hashRate, now); + } catch (error) { + console.error(`Failed to persist SV2 client hashrate: ${error.message}`); + } + } + + private getHashRatePersistIntervalMs(): number { + const configured = Number(this.configService.get('CLIENT_HASHRATE_PERSIST_INTERVAL_MS')); + if (Number.isFinite(configured) && configured >= 0) { + return configured; + } + + return DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS; + } + private parseUserIdentity(userIdentity: string): { address: string; workerName: string } { const parts = userIdentity.split('.'); const rawAddress = parts[0] ?? ''; diff --git a/src/services/datum.service.spec.ts b/src/services/datum.service.spec.ts index 384b5cc..338b585 100644 --- a/src/services/datum.service.spec.ts +++ b/src/services/datum.service.spec.ts @@ -179,6 +179,7 @@ describe('DatumService job validation', () => { jest.useFakeTimers().setSystemTime(new Date('2026-06-13T15:00:00.000Z')); const clientService = { updateBestDifficultyIfHigher: jest.fn().mockResolvedValue(undefined), + updateHashRate: jest.fn().mockResolvedValue(undefined), }; const redisMessagingService = { setClientPresence: jest.fn().mockResolvedValue(undefined), @@ -197,6 +198,7 @@ describe('DatumService job validation', () => { hashRate: 0, }, statistics: new StratumV1ClientStatistics(1), + lastHashRatePersistedAt: 0, }; await service.updateAcceptedSharePresence(state, 'tb1qdatum', 'datum-worker', 10, 1); @@ -220,6 +222,11 @@ describe('DatumService job validation', () => { const lastPresence = redisMessagingService.setClientPresence.mock.calls.at(-1)[0]; expect(lastPresence.hashRate).toBeGreaterThan(0); expect(state.clientEntity.hashRate).toBe(lastPresence.hashRate); + expect(clientService.updateHashRate).toHaveBeenCalledWith( + '3db0db03-3a62-4e3b-91bc-243adff4b542', + lastPresence.hashRate, + new Date('2026-06-13T15:01:02.000Z'), + ); }); it('derives submitted share difficulty from DATUM target byte', () => { diff --git a/src/services/datum.service.ts b/src/services/datum.service.ts index 43272d9..49e2e6a 100644 --- a/src/services/datum.service.ts +++ b/src/services/datum.service.ts @@ -44,6 +44,7 @@ import { parsePayoutModePorts, PayoutMode } from '../types/payout-mode'; const DEFAULT_DATUM_SHARE_DIFFICULTY = 1; const DEFAULT_DATUM_PING_INTERVAL_MS = 30_000; +const DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS = 60_000; @Injectable() export class DatumService implements OnModuleInit { @@ -113,6 +114,7 @@ export class DatumService implements OnModuleInit { coinbaserPayoutContexts: new Map(), nextCoinbaserId: 1, statistics: new StratumV1ClientStatistics(this.getConfiguredDatumShareDifficulty()), + lastHashRatePersistedAt: 0, }; console.log(`[DATUM ${state.sessionId}] connection accepted from ${socket.remoteAddress}:${socket.remotePort}`); @@ -389,6 +391,7 @@ export class DatumService implements OnModuleInit { await state.statistics.addShares(state.clientEntity, creditedDifficulty); state.clientEntity.hashRate = state.statistics.hashRate; + await this.persistClientHashRate(state, new Date()); if (submissionDifficulty > Number(state.clientEntity.bestDifficulty ?? 0)) { await this.clientService.updateBestDifficultyIfHigher(state.clientEntity.id, submissionDifficulty); state.clientEntity.bestDifficulty = submissionDifficulty; @@ -774,6 +777,39 @@ export class DatumService implements OnModuleInit { } } + private async persistClientHashRate(state: DatumClientState, now: Date): Promise { + if (state.clientEntity?.id == null) { + return; + } + + const hashRate = Number(state.statistics?.hashRate ?? 0); + if (!Number.isFinite(hashRate) || hashRate <= 0) { + return; + } + + const intervalMs = this.getHashRatePersistIntervalMs(); + const nowMs = now.getTime(); + if (state.lastHashRatePersistedAt > 0 && nowMs - state.lastHashRatePersistedAt < intervalMs) { + return; + } + + state.lastHashRatePersistedAt = nowMs; + try { + await this.clientService.updateHashRate(state.clientEntity.id, hashRate, now); + } catch (error) { + console.error(`Failed to persist DATUM client hashrate: ${error.message}`); + } + } + + private getHashRatePersistIntervalMs(): number { + const configured = Number(this.configService.get('CLIENT_HASHRATE_PERSIST_INTERVAL_MS')); + if (Number.isFinite(configured) && configured >= 0) { + return configured; + } + + return DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS; + } + private async destroyClient(state: DatumClientState): Promise { this.stopDatumPing(state); if (state.clientEntity?.id == null) { @@ -937,6 +973,7 @@ interface DatumClientState { coinbaserPayoutContexts: Map; nextCoinbaserId: number; statistics: StratumV1ClientStatistics; + lastHashRatePersistedAt: number; coinbaseMismatchLogged?: boolean; }