Restore DB-backed connected miner hashrate

This commit is contained in:
Ben
2026-06-20 15:27:33 -04:00
parent c3ad01a38f
commit cd0cb9bfe0
7 changed files with 124 additions and 0 deletions
+1
View File
@@ -100,6 +100,7 @@ SHARE_ACCOUNTING_MAX_QUEUE_SIZE=50000
SHARE_ACCOUNTING_SUMMARY_CACHE_MS=2500 SHARE_ACCOUNTING_SUMMARY_CACHE_MS=2500
SHARE_ACCOUNTING_SUMMARY_CACHE_MAX=10000 SHARE_ACCOUNTING_SUMMARY_CACHE_MAX=10000
CLIENT_PRESENCE_ENABLED=true CLIENT_PRESENCE_ENABLED=true
CLIENT_HASHRATE_PERSIST_INTERVAL_MS=60000
SHARE_ROLLUP_ENABLED=true SHARE_ROLLUP_ENABLED=true
SHARE_ROLLUP_INTERVAL_MS=60000 SHARE_ROLLUP_INTERVAL_MS=60000
SHARE_ROLLUP_SAFETY_LAG_SECONDS=30 SHARE_ROLLUP_SAFETY_LAG_SECONDS=30
@@ -62,6 +62,9 @@ export class UserAgentReportService {
private async buildLiveReport() { private async buildLiveReport() {
const presences = await this.redisMessagingService.getAllClientPresence(); const presences = await this.redisMessagingService.getAllClientPresence();
const activePresences = await this.filterActivePresences(presences); const activePresences = await this.filterActivePresences(presences);
if (activePresences.length === 0) {
return await this.userAgentReport.find();
}
const rows = new Map<string, { const rows = new Map<string, {
userAgent: string; userAgent: string;
count: number; count: number;
+4
View File
@@ -58,6 +58,10 @@ export class ClientService {
.execute(); .execute();
} }
public async updateHashRate(id: string, hashRate: number, updatedAt = new Date()) {
return await this.clientRepository.update({ id }, { hashRate, updatedAt });
}
public async connectedClientCount(): Promise<number> { public async connectedClientCount(): Promise<number> {
return await this.clientRepository.count(); return await this.clientRepository.count();
} }
+36
View File
@@ -37,6 +37,7 @@ const TRUE_DIFF_ONE = 2.695953529101131e67;
const BLOCKED_USER_AGENT_LOG_INTERVAL_MS = 60 * 1000; const BLOCKED_USER_AGENT_LOG_INTERVAL_MS = 60 * 1000;
const VALIDATION_ERROR_LOG_INTERVAL_MS = 60 * 1000; const VALIDATION_ERROR_LOG_INTERVAL_MS = 60 * 1000;
const DEFAULT_MIN_DIFFICULTY = 0.001; const DEFAULT_MIN_DIFFICULTY = 0.001;
const DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS = 60 * 1000;
export class StratumV1Client { export class StratumV1Client {
private static blockedUserAgentLogState = new Map<string, { nextLogAt: number, suppressed: number }>(); private static blockedUserAgentLogState = new Map<string, { nextLogAt: number, suppressed: number }>();
@@ -68,6 +69,7 @@ export class StratumV1Client {
private buffer: string = ''; private buffer: string = '';
private connectionClosed = false; private connectionClosed = false;
private lastSentMiningJobTimestamp: number = null; private lastSentMiningJobTimestamp: number = null;
private lastHashRatePersistedAt = 0;
private miningSubmissionHashes = new Set<string>() private miningSubmissionHashes = new Set<string>()
@@ -704,6 +706,7 @@ export class StratumV1Client {
const now = new Date(); const now = new Date();
this.clientEntity.updatedAt = now; this.clientEntity.updatedAt = now;
this.clientEntity.hashRate = this.statistics.hashRate; this.clientEntity.hashRate = this.statistics.hashRate;
await this.persistClientHashRate(now);
await this.updateClientPresence(now); await this.updateClientPresence(now);
} catch (e) { } catch (e) {
@@ -852,6 +855,39 @@ export class StratumV1Client {
}); });
} }
private async persistClientHashRate(now: Date): Promise<void> {
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<string>('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 { private getValidationErrorSignature(errors: ValidationError[]): string {
if (errors.length === 0) { if (errors.length === 0) {
return 'unknown'; return 'unknown';
+36
View File
@@ -70,6 +70,7 @@ const DEFAULT_START_DIFFICULTY = 100000;
const DEFAULT_MIN_DIFFICULTY = 0.001; const DEFAULT_MIN_DIFFICULTY = 0.001;
const DEFAULT_TARGET_SHARES_PER_MINUTE = 2; const DEFAULT_TARGET_SHARES_PER_MINUTE = 2;
const DEFAULT_DIFFICULTY_CHECK_INTERVAL_MS = 60 * 1000; const DEFAULT_DIFFICULTY_CHECK_INTERVAL_MS = 60 * 1000;
const DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS = 60 * 1000;
const FIXED_STANDARD_EXTRANONCE2 = '0000000000000000'; const FIXED_STANDARD_EXTRANONCE2 = '0000000000000000';
const RETIRED_EXTENDED_JOB_RETENTION_MS = 5 * 60 * 1000; const RETIRED_EXTENDED_JOB_RETENTION_MS = 5 * 60 * 1000;
const SV2_AUTH_FAILURE_LOG_INTERVAL_MS = 60 * 1000; const SV2_AUTH_FAILURE_LOG_INTERVAL_MS = 60 * 1000;
@@ -133,6 +134,7 @@ export class StratumV2Client {
private creatingEntity: Promise<void> = null; private creatingEntity: Promise<void> = null;
private readonly firstChunkSummary: string; private readonly firstChunkSummary: string;
private workSelectionEnabled = false; private workSelectionEnabled = false;
private lastHashRatePersistedAt = 0;
constructor( constructor(
private readonly socket: Socket, private readonly socket: Socket,
@@ -934,6 +936,7 @@ export class StratumV2Client {
const now = new Date(); const now = new Date();
this.clientEntity.updatedAt = now; this.clientEntity.updatedAt = now;
this.clientEntity.hashRate = this.statistics.hashRate; this.clientEntity.hashRate = this.statistics.hashRate;
await this.persistClientHashRate(now);
await this.updateClientPresence(now); await this.updateClientPresence(now);
if (submissionDifficulty > this.clientEntity.bestDifficulty) { if (submissionDifficulty > this.clientEntity.bestDifficulty) {
@@ -1343,6 +1346,39 @@ export class StratumV2Client {
} }
} }
private async persistClientHashRate(now: Date): Promise<void> {
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<string>('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 } { private parseUserIdentity(userIdentity: string): { address: string; workerName: string } {
const parts = userIdentity.split('.'); const parts = userIdentity.split('.');
const rawAddress = parts[0] ?? ''; const rawAddress = parts[0] ?? '';
+7
View File
@@ -179,6 +179,7 @@ describe('DatumService job validation', () => {
jest.useFakeTimers().setSystemTime(new Date('2026-06-13T15:00:00.000Z')); jest.useFakeTimers().setSystemTime(new Date('2026-06-13T15:00:00.000Z'));
const clientService = { const clientService = {
updateBestDifficultyIfHigher: jest.fn().mockResolvedValue(undefined), updateBestDifficultyIfHigher: jest.fn().mockResolvedValue(undefined),
updateHashRate: jest.fn().mockResolvedValue(undefined),
}; };
const redisMessagingService = { const redisMessagingService = {
setClientPresence: jest.fn().mockResolvedValue(undefined), setClientPresence: jest.fn().mockResolvedValue(undefined),
@@ -197,6 +198,7 @@ describe('DatumService job validation', () => {
hashRate: 0, hashRate: 0,
}, },
statistics: new StratumV1ClientStatistics(1), statistics: new StratumV1ClientStatistics(1),
lastHashRatePersistedAt: 0,
}; };
await service.updateAcceptedSharePresence(state, 'tb1qdatum', 'datum-worker', 10, 1); 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]; const lastPresence = redisMessagingService.setClientPresence.mock.calls.at(-1)[0];
expect(lastPresence.hashRate).toBeGreaterThan(0); expect(lastPresence.hashRate).toBeGreaterThan(0);
expect(state.clientEntity.hashRate).toBe(lastPresence.hashRate); 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', () => { it('derives submitted share difficulty from DATUM target byte', () => {
+37
View File
@@ -44,6 +44,7 @@ import { parsePayoutModePorts, PayoutMode } from '../types/payout-mode';
const DEFAULT_DATUM_SHARE_DIFFICULTY = 1; const DEFAULT_DATUM_SHARE_DIFFICULTY = 1;
const DEFAULT_DATUM_PING_INTERVAL_MS = 30_000; const DEFAULT_DATUM_PING_INTERVAL_MS = 30_000;
const DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS = 60_000;
@Injectable() @Injectable()
export class DatumService implements OnModuleInit { export class DatumService implements OnModuleInit {
@@ -113,6 +114,7 @@ export class DatumService implements OnModuleInit {
coinbaserPayoutContexts: new Map(), coinbaserPayoutContexts: new Map(),
nextCoinbaserId: 1, nextCoinbaserId: 1,
statistics: new StratumV1ClientStatistics(this.getConfiguredDatumShareDifficulty()), statistics: new StratumV1ClientStatistics(this.getConfiguredDatumShareDifficulty()),
lastHashRatePersistedAt: 0,
}; };
console.log(`[DATUM ${state.sessionId}] connection accepted from ${socket.remoteAddress}:${socket.remotePort}`); 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); await state.statistics.addShares(state.clientEntity, creditedDifficulty);
state.clientEntity.hashRate = state.statistics.hashRate; state.clientEntity.hashRate = state.statistics.hashRate;
await this.persistClientHashRate(state, new Date());
if (submissionDifficulty > Number(state.clientEntity.bestDifficulty ?? 0)) { if (submissionDifficulty > Number(state.clientEntity.bestDifficulty ?? 0)) {
await this.clientService.updateBestDifficultyIfHigher(state.clientEntity.id, submissionDifficulty); await this.clientService.updateBestDifficultyIfHigher(state.clientEntity.id, submissionDifficulty);
state.clientEntity.bestDifficulty = submissionDifficulty; state.clientEntity.bestDifficulty = submissionDifficulty;
@@ -774,6 +777,39 @@ export class DatumService implements OnModuleInit {
} }
} }
private async persistClientHashRate(state: DatumClientState, now: Date): Promise<void> {
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<string>('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<void> { private async destroyClient(state: DatumClientState): Promise<void> {
this.stopDatumPing(state); this.stopDatumPing(state);
if (state.clientEntity?.id == null) { if (state.clientEntity?.id == null) {
@@ -937,6 +973,7 @@ interface DatumClientState {
coinbaserPayoutContexts: Map<number, DatumCoinbaserPayoutContext>; coinbaserPayoutContexts: Map<number, DatumCoinbaserPayoutContext>;
nextCoinbaserId: number; nextCoinbaserId: number;
statistics: StratumV1ClientStatistics; statistics: StratumV1ClientStatistics;
lastHashRatePersistedAt: number;
coinbaseMismatchLogged?: boolean; coinbaseMismatchLogged?: boolean;
} }