From d1fb186e5dddea146ebb322487d6c26a73453da4 Mon Sep 17 00:00:00 2001 From: Ben Date: Sat, 20 Jun 2026 16:03:34 -0400 Subject: [PATCH] Remove Redis client presence --- .env.example | 1 - docker-compose.external-db.yml | 1 - docker-compose.yml | 1 - full-setup/docker-compose-mainnet.yml | 1 - .../user-agent-report.service.spec.ts | 106 ++++----- .../user-agent-report.service.ts | 78 ++----- src/ORM/client/client.service.ts | 5 +- .../client/client.controller.spec.ts | 45 ++-- src/controllers/client/client.controller.ts | 26 ++- src/models/StratumV1Client.spec.ts | 20 +- src/models/StratumV1Client.ts | 28 --- src/models/StratumV2Client.spec.ts | 18 +- src/models/StratumV2Client.ts | 27 --- src/services/datum.service.spec.ts | 21 +- src/services/datum.service.ts | 41 ---- src/services/redis-messaging.service.spec.ts | 74 ------- src/services/redis-messaging.service.ts | 209 +----------------- src/services/stratum-v1.service.spec.ts | 5 +- src/services/stratum-v1.service.ts | 1 - test/timescale-redis.integration-spec.ts | 16 +- 20 files changed, 111 insertions(+), 613 deletions(-) diff --git a/.env.example b/.env.example index 431003c..f69b7a4 100644 --- a/.env.example +++ b/.env.example @@ -99,7 +99,6 @@ SHARE_ACCOUNTING_FLUSH_INTERVAL_MS=25 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 diff --git a/docker-compose.external-db.yml b/docker-compose.external-db.yml index 5f0119a..48840cb 100644 --- a/docker-compose.external-db.yml +++ b/docker-compose.external-db.yml @@ -100,7 +100,6 @@ services: SHARE_ACCOUNTING_MAX_QUEUE_SIZE: ${SHARE_ACCOUNTING_MAX_QUEUE_SIZE:-50000} SHARE_ACCOUNTING_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MS:-2500} SHARE_ACCOUNTING_SUMMARY_CACHE_MAX: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MAX:-10000} - CLIENT_PRESENCE_ENABLED: ${CLIENT_PRESENCE_ENABLED:-true} SHARE_ROLLUP_ENABLED: ${SHARE_ROLLUP_ENABLED:-true} SHARE_ROLLUP_INTERVAL_MS: ${SHARE_ROLLUP_INTERVAL_MS:-60000} SHARE_ROLLUP_SAFETY_LAG_SECONDS: ${SHARE_ROLLUP_SAFETY_LAG_SECONDS:-30} diff --git a/docker-compose.yml b/docker-compose.yml index b9fc99e..072efb1 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -124,7 +124,6 @@ services: PPLNS_SV2_TDP_PORTS: ${PPLNS_SV2_TDP_PORTS:-} DATUM_PORTS: ${DATUM_PORTS:-} PPLNS_DATUM_PORTS: ${PPLNS_DATUM_PORTS:-} - CLIENT_PRESENCE_ENABLED: ${CLIENT_PRESENCE_ENABLED:-true} healthcheck: test: ["CMD-SHELL", "node -e \"const http=require('http'); const https=require('https'); const secure=process.env.API_SECURE==='true'; const client=secure?https:http; const req=client.get({hostname:'127.0.0.1',port:process.env.API_PORT||3334,path:'/api/network',rejectUnauthorized:false},res=>process.exit(res.statusCode<500?0:1)); req.on('error',()=>process.exit(1)); req.setTimeout(5000,()=>{req.destroy(); process.exit(1);});\""] interval: 30s diff --git a/full-setup/docker-compose-mainnet.yml b/full-setup/docker-compose-mainnet.yml index fe4ab59..cf0726d 100644 --- a/full-setup/docker-compose-mainnet.yml +++ b/full-setup/docker-compose-mainnet.yml @@ -113,7 +113,6 @@ services: 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} - CLIENT_PRESENCE_ENABLED: ${CLIENT_PRESENCE_ENABLED:-true} SHARE_ROLLUP_ENABLED: ${SHARE_ROLLUP_ENABLED:-true} SHARE_ROLLUP_INTERVAL_MS: ${SHARE_ROLLUP_INTERVAL_MS:-60000} SHARE_ROLLUP_SAFETY_LAG_SECONDS: ${SHARE_ROLLUP_SAFETY_LAG_SECONDS:-30} diff --git a/src/ORM/_views/user-agent-report/user-agent-report.service.spec.ts b/src/ORM/_views/user-agent-report/user-agent-report.service.spec.ts index 245ac4e..5eeca58 100644 --- a/src/ORM/_views/user-agent-report/user-agent-report.service.spec.ts +++ b/src/ORM/_views/user-agent-report/user-agent-report.service.spec.ts @@ -11,7 +11,7 @@ describe('UserAgentReportService', () => { } }); - it('uses Redis live presence in API-only processes when workers are connected', async () => { + it('uses active database clients for live API-only reports', async () => { process.env.API_ONLY = 'true'; const userAgentReport = { find: jest.fn().mockResolvedValue([ @@ -23,26 +23,17 @@ describe('UserAgentReportService', () => { }, ]), }; - const clientRepository = { - find: jest.fn().mockResolvedValue([{ id: 'client-1' }]), - }; + const clientRepository = createClientRepository([ + { + userAgent: 'unknown/sv2', + count: '1', + bestDifficulty: 456, + totalHashRate: '123', + }, + ]); const redisMessagingService = { getJsonCache: jest.fn().mockResolvedValue(null), - getAllClientPresence: jest.fn().mockResolvedValue([ - { - clientId: 'client-1', - address: 'tb1qworker', - clientName: 'sv2gateway', - sessionId: 'session', - userAgent: 'unknown/sv2', - startTime: '2026-01-01T00:00:00.000Z', - lastSeen: '2026-01-01T00:00:10.000Z', - hashRate: 123, - bestDifficulty: 456, - }, - ]), setJsonCache: jest.fn().mockResolvedValue(undefined), - removeClientPresence: jest.fn().mockResolvedValue(undefined), }; const service = new UserAgentReportService( userAgentReport as any, @@ -59,9 +50,10 @@ describe('UserAgentReportService', () => { }, ]); expect(userAgentReport.find).not.toHaveBeenCalled(); + expect(clientRepository.createQueryBuilder).toHaveBeenCalled(); }); - it('falls back to the materialized view in API-only processes when live presence is empty', async () => { + it('falls back to the materialized view when no active database clients exist', async () => { process.env.API_ONLY = 'true'; const viewRows = [ { @@ -74,14 +66,10 @@ describe('UserAgentReportService', () => { const userAgentReport = { find: jest.fn().mockResolvedValue(viewRows), }; - const clientRepository = { - find: jest.fn().mockResolvedValue([]), - }; + const clientRepository = createClientRepository([]); const redisMessagingService = { getJsonCache: jest.fn().mockResolvedValue([]), - getAllClientPresence: jest.fn().mockResolvedValue([]), setJsonCache: jest.fn().mockResolvedValue(undefined), - removeClientPresence: jest.fn().mockResolvedValue(undefined), }; const service = new UserAgentReportService( userAgentReport as any, @@ -92,42 +80,23 @@ describe('UserAgentReportService', () => { await expect(service.getReport()).resolves.toEqual(viewRows); }); - it('filters and removes stale Redis presence rows with no active database client', async () => { + it('returns a non-empty cached live report before querying the database', async () => { process.env.API_ONLY = 'true'; + const cachedRows = [ + { + userAgent: 'cached', + count: '2', + bestDifficulty: 8, + totalHashRate: '16', + }, + ]; const userAgentReport = { - find: jest.fn().mockResolvedValue([]), - }; - const clientRepository = { - find: jest.fn().mockResolvedValue([{ id: 'active-client' }]), + find: jest.fn(), }; + const clientRepository = createClientRepository([]); const redisMessagingService = { - getJsonCache: jest.fn().mockResolvedValue(null), - getAllClientPresence: jest.fn().mockResolvedValue([ - { - clientId: 'active-client', - address: 'tb1qactive', - clientName: 'worker', - sessionId: 'active', - userAgent: 'bitaxe', - startTime: '2026-01-01T00:00:00.000Z', - lastSeen: '2026-01-01T00:00:10.000Z', - hashRate: 100, - bestDifficulty: 200, - }, - { - clientId: 'stale-client', - address: 'tb1qstale', - clientName: 'worker', - sessionId: 'stale', - userAgent: 'bitaxe', - startTime: '2026-01-01T00:00:00.000Z', - lastSeen: '2026-01-01T00:00:10.000Z', - hashRate: 300, - bestDifficulty: 400, - }, - ]), + getJsonCache: jest.fn().mockResolvedValue(cachedRows), setJsonCache: jest.fn().mockResolvedValue(undefined), - removeClientPresence: jest.fn().mockResolvedValue(undefined), }; const service = new UserAgentReportService( userAgentReport as any, @@ -135,15 +104,22 @@ describe('UserAgentReportService', () => { redisMessagingService as any, ); - await expect(service.getReport()).resolves.toEqual([ - { - userAgent: 'bitaxe', - count: '1', - bestDifficulty: 200, - totalHashRate: '100', - }, - ]); - expect(redisMessagingService.removeClientPresence) - .toHaveBeenCalledWith('stale-client', 'tb1qstale'); + await expect(service.getReport()).resolves.toEqual(cachedRows); + expect(clientRepository.createQueryBuilder).not.toHaveBeenCalled(); }); }); + +function createClientRepository(rows: unknown[]) { + const queryBuilder = { + select: jest.fn().mockReturnThis(), + addSelect: jest.fn().mockReturnThis(), + where: jest.fn().mockReturnThis(), + groupBy: jest.fn().mockReturnThis(), + orderBy: jest.fn().mockReturnThis(), + getRawMany: jest.fn().mockResolvedValue(rows), + }; + + return { + createQueryBuilder: jest.fn().mockReturnValue(queryBuilder), + }; +} 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 ce42c88..7a88662 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 @@ -1,6 +1,6 @@ import { Injectable } from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; -import { In, Repository } from 'typeorm'; +import { Repository } from 'typeorm'; import { ClientEntity } from '../../client/client.entity'; import { UserAgentReportView } from './user-agent-report.view'; @@ -60,76 +60,28 @@ export class UserAgentReportService { } private async buildLiveReport() { - const presences = await this.redisMessagingService.getAllClientPresence(); - const activePresences = await this.filterActivePresences(presences); - if (activePresences.length === 0) { + const rows = await this.clientRepository + .createQueryBuilder('client') + .select('COALESCE(NULLIF(client.userAgent, \'\'), \'Other\')', 'userAgent') + .addSelect('COUNT(*)', 'count') + .addSelect('MAX(client.bestDifficulty)', 'bestDifficulty') + .addSelect('COALESCE(SUM(client.hashRate), 0)', 'totalHashRate') + .where('client.deletedAt IS NULL') + .groupBy('COALESCE(NULLIF(client.userAgent, \'\'), \'Other\')') + .orderBy('"totalHashRate"', 'DESC') + .getRawMany(); + + if (rows.length === 0) { return await this.userAgentReport.find(); } - const rows = new Map(); - - activePresences.forEach(presence => { - const userAgent = presence.userAgent == null || presence.userAgent.length === 0 - ? 'Other' - : presence.userAgent; - const row = rows.get(userAgent) ?? { - userAgent, - count: 0, - bestDifficulty: 0, - totalHashRate: 0, - }; - row.count++; - row.bestDifficulty = Math.max( - row.bestDifficulty, - Number(presence.bestDifficulty ?? 0), - ); - row.totalHashRate += Number(presence.hashRate ?? 0); - rows.set(userAgent, row); - }); - - const report = [...rows.values()] - .sort((left, right) => right.totalHashRate - left.totalHashRate) - .map(row => ({ - userAgent: row.userAgent, - count: row.count.toString(), - bestDifficulty: row.bestDifficulty, - totalHashRate: row.totalHashRate.toString(), - })); await this.redisMessagingService - .setJsonCache(this.liveReportCacheKey, report, 60 * 1000) + .setJsonCache(this.liveReportCacheKey, rows, 60 * 1000) .catch(error => { console.error(`Live user-agent report cache write failed: ${error.message}`); }); - return report; - } - - private async filterActivePresences(presences: Awaited>) { - if (presences.length === 0) { - return []; - } - - const activeClients = await this.clientRepository.find({ - select: { - id: true, - }, - where: { - id: In(presences.map(presence => presence.clientId)), - }, - }); - const activeIds = new Set(activeClients.map(client => client.id)); - const stalePresences = presences.filter(presence => !activeIds.has(presence.clientId)); - void Promise.all(stalePresences.map(presence => { - return this.redisMessagingService.removeClientPresence(presence.clientId, presence.address) - .catch(() => undefined); - })); - - return presences.filter(presence => activeIds.has(presence.clientId)); + return rows; } public async refreshReport() { diff --git a/src/ORM/client/client.service.ts b/src/ORM/client/client.service.ts index f81b40b..255fca1 100644 --- a/src/ORM/client/client.service.ts +++ b/src/ORM/client/client.service.ts @@ -87,7 +87,10 @@ export class ClientService { return await this.clientRepository.find({ where: { address - } + }, + order: { + updatedAt: 'DESC', + }, }) } diff --git a/src/controllers/client/client.controller.spec.ts b/src/controllers/client/client.controller.spec.ts index 3fbc341..4cc1b2f 100644 --- a/src/controllers/client/client.controller.spec.ts +++ b/src/controllers/client/client.controller.spec.ts @@ -5,12 +5,11 @@ import { ClientStatisticsService } from '../../ORM/client-statistics/client-stat import { ClientService } from '../../ORM/client/client.service'; import { PayoutSnapshotService } from '../../ORM/payout-snapshot/payout-snapshot.service'; import { ShareAccountingService } from '../../ORM/share-accounting/share-accounting.service'; -import { RedisMessagingService } from '../../services/redis-messaging.service'; import { ClientController } from './client.controller'; describe('ClientController', () => { let controller: ClientController; - let clientService: { getBySessionId: jest.Mock; getActiveIds: jest.Mock }; + let clientService: { getByAddress: jest.Mock; getBySessionId: jest.Mock }; let addressSettingsService: { getSettings: jest.Mock }; let shareAccountingService: { getAddressSummary: jest.Mock; @@ -19,7 +18,6 @@ describe('ClientController', () => { getSessionSummaries: jest.Mock; }; let payoutSnapshotService: { getLatestExpectedPayoutForAddress: jest.Mock }; - let redisMessagingService: { getClientPresenceByAddress: jest.Mock; removeClientPresence: jest.Mock }; beforeEach(async () => { const module: TestingModule = await Test.createTestingModule({ @@ -28,10 +26,9 @@ describe('ClientController', () => { { provide: ClientService, useValue: { - getByAddress: jest.fn(), + getByAddress: jest.fn().mockResolvedValue([]), getByName: jest.fn(), getBySessionId: jest.fn(), - getActiveIds: jest.fn(async (ids: string[]) => new Set(ids)), }, }, { @@ -63,13 +60,6 @@ describe('ClientController', () => { getLatestExpectedPayoutForAddress: jest.fn().mockResolvedValue(null), }, }, - { - provide: RedisMessagingService, - useValue: { - getClientPresenceByAddress: jest.fn().mockResolvedValue([]), - removeClientPresence: jest.fn().mockResolvedValue(undefined), - }, - }, ], }).compile(); @@ -79,7 +69,6 @@ describe('ClientController', () => { addressSettingsService = module.get(AddressSettingsService); shareAccountingService = module.get(ShareAccountingService); payoutSnapshotService = module.get(PayoutSnapshotService); - redisMessagingService = module.get(RedisMessagingService); }); it('should be defined', () => { @@ -87,9 +76,9 @@ describe('ClientController', () => { }); it('should expose the existing address best difficulty in accounting when rollup best is empty', async () => { - redisMessagingService.getClientPresenceByAddress.mockResolvedValue([ + clientService.getByAddress.mockResolvedValue([ { - clientId: '92f5302f-5e32-487e-af67-f56fd78b13c7', + id: '92f5302f-5e32-487e-af67-f56fd78b13c7', sessionId: 'abcd1234', clientName: 'worker', bestDifficulty: 64, @@ -97,6 +86,7 @@ describe('ClientController', () => { startTime: '2026-06-08T12:00:00.000Z', lastSeen: '2026-06-08T12:10:00.000Z', address: 'bc1qtest', + payoutMode: 'pplns', }, ]); addressSettingsService.getSettings.mockResolvedValue({ bestDifficulty: 4096 }); @@ -154,30 +144,31 @@ describe('ClientController', () => { expect(payoutSnapshotService.getLatestExpectedPayoutForAddress).not.toHaveBeenCalled(); }); - it('should hide stale Redis workers that are no longer active in the database', async () => { - redisMessagingService.getClientPresenceByAddress.mockResolvedValue([ + it('should expose active database workers for an address', async () => { + clientService.getByAddress.mockResolvedValue([ { - clientId: 'active-client', + id: 'active-client', address: 'bc1qtest', sessionId: 'active1', clientName: 'active-worker', + payoutMode: 'pplns', bestDifficulty: 64, hashRate: 1024, startTime: '2026-06-08T12:00:00.000Z', - lastSeen: '2026-06-08T12:10:00.000Z', + updatedAt: '2026-06-08T12:10:00.000Z', }, { - clientId: 'stale-client', + id: 'solo-client', address: 'bc1qtest', - sessionId: 'stale1', - clientName: 'stale-worker', + sessionId: 'solo1', + clientName: 'solo-worker', + payoutMode: 'solo', bestDifficulty: 128, hashRate: 2048, startTime: '2026-06-08T12:00:00.000Z', - lastSeen: '2026-06-08T12:10:00.000Z', + updatedAt: '2026-06-08T12:10:00.000Z', }, ]); - clientService.getActiveIds.mockResolvedValue(new Set(['active-client'])); addressSettingsService.getSettings.mockResolvedValue(null); shareAccountingService.getSessionSummaries.mockResolvedValue(new Map()); shareAccountingService.getAddressSummary.mockResolvedValue({ @@ -186,16 +177,16 @@ describe('ClientController', () => { bestSubmissionDifficulty: 0, }); - await expect(controller.getClientInfo('bc1qtest')).resolves.toMatchObject({ + await expect(controller.getClientInfo('bc1qtest', 'pplns')).resolves.toMatchObject({ workersCount: 1, workers: [ { sessionId: 'active1', name: 'active-worker', + payoutMode: 'pplns', + hashRate: 1024, }, ], }); - expect(redisMessagingService.removeClientPresence) - .toHaveBeenCalledWith('stale-client', 'bc1qtest'); }); }); diff --git a/src/controllers/client/client.controller.ts b/src/controllers/client/client.controller.ts index c413dee..5444484 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 { PayoutSnapshotService } from '../../ORM/payout-snapshot/payout-snapshot.service'; import { ShareAccountingService } from '../../ORM/share-accounting/share-accounting.service'; -import { RedisMessagingService } from '../../services/redis-messaging.service'; import { normalizePayoutMode, PayoutMode } from '../../types/payout-mode'; @@ -18,7 +17,6 @@ export class ClientController { private readonly addressSettingsService: AddressSettingsService, private readonly shareAccountingService: ShareAccountingService, private readonly payoutSnapshotService: PayoutSnapshotService, - private readonly redisMessagingService: RedisMessagingService ) { } @@ -179,14 +177,20 @@ export class ClientController { } private async getActiveAddressWorkers(address: string, payoutMode?: PayoutMode) { - const workers = await this.redisMessagingService.getClientPresenceByAddress(address); - const activeIds = await this.clientService.getActiveIds(workers.map(worker => worker.clientId)); - const staleWorkers = workers.filter(worker => !activeIds.has(worker.clientId)); - void Promise.all(staleWorkers.map(worker => { - return this.redisMessagingService.removeClientPresence(worker.clientId, worker.address) - .catch(() => undefined); - })); - - return workers.filter(worker => activeIds.has(worker.clientId) && (payoutMode == null || worker.payoutMode === payoutMode)); + const workers = await this.clientService.getByAddress(address); + return workers + .filter(worker => payoutMode == null || worker.payoutMode === payoutMode) + .map(worker => ({ + clientId: worker.id, + address: worker.address, + clientName: worker.clientName, + sessionId: worker.sessionId, + payoutMode: worker.payoutMode, + userAgent: worker.userAgent, + startTime: worker.startTime, + lastSeen: worker.updatedAt, + hashRate: Number(worker.hashRate ?? 0), + bestDifficulty: Number(worker.bestDifficulty ?? 0), + })); } } diff --git a/src/models/StratumV1Client.spec.ts b/src/models/StratumV1Client.spec.ts index c99a840..7c48e98 100644 --- a/src/models/StratumV1Client.spec.ts +++ b/src/models/StratumV1Client.spec.ts @@ -41,7 +41,7 @@ describe('StratumV1Client', () => { let configService: ConfigService; let addressSettings: AddressSettingsService; let shareAccountingService: { recordAcceptedShare: jest.Mock }; - let redisMessagingService: { setClientPresence: jest.Mock; removeClientPresence: jest.Mock }; + let redisMessagingService: Record; let client: StratumV1Client; @@ -80,6 +80,7 @@ describe('StratumV1Client', () => { }), connectedClientCount: jest.fn(async () => clients.size), updateBestDifficultyIfHigher: jest.fn().mockResolvedValue({ affected: 1 }), + updateHashRate: jest.fn().mockResolvedValue(undefined), } as any; configService = { @@ -132,10 +133,7 @@ describe('StratumV1Client', () => { shareAccountingService = { recordAcceptedShare: jest.fn().mockResolvedValue(undefined), }; - redisMessagingService = { - setClientPresence: jest.fn().mockResolvedValue(undefined), - removeClientPresence: jest.fn().mockResolvedValue(undefined), - }; + redisMessagingService = {}; client = new StratumV1Client( @@ -185,11 +183,6 @@ describe('StratumV1Client', () => { await Promise.all([client.destroy(), client.destroy()]); - expect(redisMessagingService.removeClientPresence).toHaveBeenCalledTimes(1); - expect(redisMessagingService.removeClientPresence).toHaveBeenCalledWith( - '00000000-0000-4000-8000-000000000001', - 'tb1qcleanup', - ); expect(clientService.delete).toHaveBeenCalledTimes(1); expect(clientService.delete).toHaveBeenCalledWith('00000000-0000-4000-8000-000000000001'); expect(unsubscribe).toHaveBeenCalledTimes(1); @@ -415,13 +408,6 @@ describe('StratumV1Client', () => { creditedDifficulty: 0, isBlockCandidate: false, })); - expect(redisMessagingService.setClientPresence).toHaveBeenCalledWith(expect.objectContaining({ - address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', - clientName: 'bitaxe3', - sessionId: MockRecording1.EXTRA_NONCE, - })); - - }); it('should use the header-only fast path for non-block submissions', async () => { diff --git a/src/models/StratumV1Client.ts b/src/models/StratumV1Client.ts index a5001e1..4987c1e 100644 --- a/src/models/StratumV1Client.ts +++ b/src/models/StratumV1Client.ts @@ -124,8 +124,6 @@ export class StratumV1Client { if (this.clientEntity?.id) { const clientId = this.clientEntity.id; - const address = this.clientEntity.address; - await this.redisMessagingService?.removeClientPresence(clientId, address); await this.clientService.delete(clientId); } } @@ -527,36 +525,12 @@ export class StratumV1Client { payoutMode: this.payoutMode, bestDifficulty: 0 }); - await this.updateClientPresence(new Date()); })(); } await this.creatingEntity; } - private async updateClientPresence(lastSeen: Date): Promise { - if (this.clientEntity == null) { - return; - } - - try { - await this.redisMessagingService?.setClientPresence({ - clientId: this.clientEntity.id, - address: this.clientEntity.address, - clientName: this.clientEntity.clientName, - sessionId: this.clientEntity.sessionId, - payoutMode: this.payoutMode, - userAgent: this.clientEntity.userAgent, - startTime: new Date(this.clientEntity.startTime).toISOString(), - lastSeen: lastSeen.toISOString(), - hashRate: this.statistics?.hashRate ?? 0, - bestDifficulty: Number(this.clientEntity.bestDifficulty ?? 0), - }); - } catch (error) { - console.error(`Failed to update SV1 client presence: ${error.message}`); - } - } - private async handleMiningSubmission(submission: MiningSubmitMessage) { const job = this.stratumV1JobsService.getJobById(submission.jobId); @@ -707,7 +681,6 @@ export class StratumV1Client { this.clientEntity.updatedAt = now; this.clientEntity.hashRate = this.statistics.hashRate; await this.persistClientHashRate(now); - await this.updateClientPresence(now); } catch (e) { console.log(e); @@ -717,7 +690,6 @@ export class StratumV1Client { await this.clientService.updateBestDifficultyIfHigher(this.clientEntity.id, submissionDifficulty); this.clientEntity.bestDifficulty = submissionDifficulty; await this.addressSettingsService.updateBestDifficultyIfHigher(this.clientAuthorization.address, submissionDifficulty, this.clientEntity.userAgent); - await this.updateClientPresence(new Date()); } diff --git a/src/models/StratumV2Client.spec.ts b/src/models/StratumV2Client.spec.ts index 9106db3..2c9b272 100644 --- a/src/models/StratumV2Client.spec.ts +++ b/src/models/StratumV2Client.spec.ts @@ -301,8 +301,8 @@ describe('StratumV2Client extended channels', () => { expect((client as any).channels.get(1).extendedJobs.has(100)).toBe(false); }); - it('records accepted SV2 shares before presence updates', async () => { - const { client, shareAccountingService, redisMessagingService, jobTemplate } = await createClient(); + it('records accepted SV2 shares and persists DB hashrate', async () => { + const { client, clientService, shareAccountingService, jobTemplate } = await createClient(); (client as any).address = 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4'; (client as any).workerName = 'worker'; (client as any).sessionId = 'sv2-session'; @@ -333,10 +333,6 @@ describe('StratumV2Client extended channels', () => { nonce: 123, extraNonce2: 'c708000000000000', })); - const accountingCallOrder = shareAccountingService.recordAcceptedShare.mock.invocationCallOrder[0]; - const postAccountingPresenceUpdates = redisMessagingService.setClientPresence.mock.invocationCallOrder - .filter(callOrder => callOrder > accountingCallOrder); - expect(postAccountingPresenceUpdates.length).toBeGreaterThan(0); }); it('does not submit a block when only the reported SV2 difficulty is huge', async () => { @@ -422,11 +418,12 @@ describe('StratumV2Client extended channels', () => { client: StratumV2Client; sentFrames: any[]; bitcoinRpcService: { SUBMIT_BLOCK: jest.Mock }; + clientService: { updateHashRate: jest.Mock }; blocksService: { save: jest.Mock }; notificationService: { notifySubscribersBlockFound: jest.Mock }; addressSettingsService: { resetBestDifficultyAndShares: jest.Mock; updateBestDifficultyIfHigher: jest.Mock }; shareAccountingService: { recordAcceptedShare: jest.Mock }; - redisMessagingService: { setClientPresence: jest.Mock; removeClientPresence: jest.Mock }; + redisMessagingService: Record; jobTemplate: any; }> { const blockTemplate$ = new BehaviorSubject(MockRecording1.BLOCK_TEMPLATE); @@ -469,6 +466,7 @@ describe('StratumV2Client extended channels', () => { insert: jest.fn().mockResolvedValue(clientEntity), delete: jest.fn().mockResolvedValue(undefined), updateBestDifficultyIfHigher: jest.fn().mockResolvedValue(undefined), + updateHashRate: jest.fn().mockResolvedValue(undefined), }; const shareAccountingService = { recordAcceptedShare: jest.fn().mockResolvedValue(undefined), @@ -479,10 +477,7 @@ describe('StratumV2Client extended channels', () => { const blocksService = { save: jest.fn().mockResolvedValue(undefined), }; - const redisMessagingService = { - setClientPresence: jest.fn().mockResolvedValue(undefined), - removeClientPresence: jest.fn().mockResolvedValue(undefined), - }; + const redisMessagingService = {}; const addressSettingsService = { resetBestDifficultyAndShares: jest.fn().mockResolvedValue(undefined), updateBestDifficultyIfHigher: jest.fn().mockResolvedValue(undefined), @@ -548,6 +543,7 @@ describe('StratumV2Client extended channels', () => { client, sentFrames, bitcoinRpcService, + clientService, blocksService, notificationService, addressSettingsService, diff --git a/src/models/StratumV2Client.ts b/src/models/StratumV2Client.ts index 1b1b9ae..f310991 100644 --- a/src/models/StratumV2Client.ts +++ b/src/models/StratumV2Client.ts @@ -191,7 +191,6 @@ export class StratumV2Client { } this.channels.clear(); if (this.clientEntity?.id != null) { - await this.redisMessagingService?.removeClientPresence(this.clientEntity.id, this.clientEntity.address); await this.clientService.delete(this.clientEntity.id); } } @@ -937,7 +936,6 @@ export class StratumV2Client { this.clientEntity.updatedAt = now; this.clientEntity.hashRate = this.statistics.hashRate; await this.persistClientHashRate(now); - await this.updateClientPresence(now); if (submissionDifficulty > this.clientEntity.bestDifficulty) { await this.clientService.updateBestDifficultyIfHigher(this.clientEntity.id, submissionDifficulty); @@ -947,7 +945,6 @@ export class StratumV2Client { submissionDifficulty, this.userAgent, ); - await this.updateClientPresence(new Date()); } } @@ -1316,36 +1313,12 @@ export class StratumV2Client { payoutMode: this.payoutMode, bestDifficulty: 0, }); - await this.updateClientPresence(new Date()); })(); } await this.creatingEntity; } - private async updateClientPresence(lastSeen: Date): Promise { - if (this.clientEntity == null) { - return; - } - - try { - await this.redisMessagingService?.setClientPresence({ - clientId: this.clientEntity.id, - address: this.clientEntity.address, - clientName: this.clientEntity.clientName, - sessionId: this.clientEntity.sessionId, - payoutMode: this.payoutMode, - userAgent: this.clientEntity.userAgent, - startTime: new Date(this.clientEntity.startTime).toISOString(), - lastSeen: lastSeen.toISOString(), - hashRate: this.statistics.hashRate, - bestDifficulty: Number(this.clientEntity.bestDifficulty ?? 0), - }); - } catch (error) { - console.error(`Failed to update SV2 client presence: ${error.message}`); - } - } - private async persistClientHashRate(now: Date): Promise { if (this.clientEntity?.id == null) { return; diff --git a/src/services/datum.service.spec.ts b/src/services/datum.service.spec.ts index 338b585..c473ac8 100644 --- a/src/services/datum.service.spec.ts +++ b/src/services/datum.service.spec.ts @@ -175,16 +175,13 @@ describe('DatumService job validation', () => { expect(socket.write).not.toHaveBeenCalled(); }); - it('publishes DATUM client presence with smoothed hashrate after accepted shares', async () => { + it('persists DATUM client hashrate after accepted shares', async () => { 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), - }; - const service = createService({ clientService, redisMessagingService }) as any; + const service = createService({ clientService }) as any; const state = { sessionId: 'datum-session', userAgent: 'datum/test', @@ -211,20 +208,10 @@ describe('DatumService job validation', () => { '3db0db03-3a62-4e3b-91bc-243adff4b542', 30, ); - expect(redisMessagingService.setClientPresence).toHaveBeenLastCalledWith(expect.objectContaining({ - clientId: '3db0db03-3a62-4e3b-91bc-243adff4b542', - address: 'tb1qdatum', - clientName: 'datum-worker', - sessionId: 'datum-session', - hashRate: expect.any(Number), - bestDifficulty: 30, - })); - const lastPresence = redisMessagingService.setClientPresence.mock.calls.at(-1)[0]; - expect(lastPresence.hashRate).toBeGreaterThan(0); - expect(state.clientEntity.hashRate).toBe(lastPresence.hashRate); + expect(state.clientEntity.hashRate).toBeGreaterThan(0); expect(clientService.updateHashRate).toHaveBeenCalledWith( '3db0db03-3a62-4e3b-91bc-243adff4b542', - lastPresence.hashRate, + state.clientEntity.hashRate, new Date('2026-06-13T15:01:02.000Z'), ); }); diff --git a/src/services/datum.service.ts b/src/services/datum.service.ts index 49e2e6a..c5d6ccc 100644 --- a/src/services/datum.service.ts +++ b/src/services/datum.service.ts @@ -396,22 +396,6 @@ export class DatumService implements OnModuleInit { await this.clientService.updateBestDifficultyIfHigher(state.clientEntity.id, submissionDifficulty); state.clientEntity.bestDifficulty = submissionDifficulty; } - - await this.redisMessagingService?.setClientPresence({ - clientId: state.clientEntity.id, - address, - clientName: workerName, - sessionId: state.sessionId, - payoutMode: state.payoutMode, - userAgent: state.userAgent, - startTime: new Date(state.clientEntity.startTime).toISOString(), - lastSeen: new Date().toISOString(), - hashRate: state.statistics.hashRate, - bestDifficulty: Math.max( - Number(state.clientEntity.bestDifficulty ?? 0), - submissionDifficulty, - ), - }); } private updateDatumJobCache(state: DatumClientState, pow: DatumPowSubmit): DatumJobCache { @@ -751,30 +735,6 @@ export class DatumService implements OnModuleInit { payoutMode: state.payoutMode, bestDifficulty: 0, }); - await this.updateClientPresence(state, new Date()); - } - - private async updateClientPresence(state: DatumClientState, lastSeen: Date): Promise { - if (state.clientEntity == null) { - return; - } - - try { - await this.redisMessagingService?.setClientPresence({ - clientId: state.clientEntity.id, - address: state.clientEntity.address, - clientName: state.clientEntity.clientName, - sessionId: state.clientEntity.sessionId, - payoutMode: state.payoutMode, - userAgent: state.clientEntity.userAgent, - startTime: new Date(state.clientEntity.startTime).toISOString(), - lastSeen: lastSeen.toISOString(), - hashRate: state.statistics.hashRate, - bestDifficulty: Number(state.clientEntity.bestDifficulty ?? 0), - }); - } catch (error) { - console.error(`Failed to update DATUM client presence: ${error.message}`); - } } private async persistClientHashRate(state: DatumClientState, now: Date): Promise { @@ -815,7 +775,6 @@ export class DatumService implements OnModuleInit { if (state.clientEntity?.id == null) { return; } - await this.redisMessagingService?.removeClientPresence(state.clientEntity.id, state.clientEntity.address); await this.clientService.delete(state.clientEntity.id); state.clientEntity = null; } diff --git a/src/services/redis-messaging.service.spec.ts b/src/services/redis-messaging.service.spec.ts index 9d4829c..6c0e7fb 100644 --- a/src/services/redis-messaging.service.spec.ts +++ b/src/services/redis-messaging.service.spec.ts @@ -61,80 +61,6 @@ describe('RedisMessagingService', () => { consoleSpy.mockRestore(); }); - it('should store, index, and remove client presence', async () => { - await service.connect(); - - await service.setClientPresence({ - clientId: '00000000-0000-4000-8000-000000000001', - address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', - clientName: 'worker', - sessionId: '57a6f098', - userAgent: 'bitaxe', - startTime: '2026-06-07T12:00:00.000Z', - lastSeen: '2026-06-07T12:01:00.000Z', - hashRate: 100, - bestDifficulty: 200, - }); - - expect(await service.getClientPresence('00000000-0000-4000-8000-000000000001')) - .toEqual(expect.objectContaining({ - clientId: '00000000-0000-4000-8000-000000000001', - address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', - hashRate: 100, - bestDifficulty: 200, - })); - expect(await service.getClientPresenceByAddress('tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4')) - .toHaveLength(1); - expect(await service.getAllClientPresence()).toHaveLength(1); - - await service.removeClientPresence('00000000-0000-4000-8000-000000000001'); - - expect(await service.getAllClientPresence()).toHaveLength(0); - }); - - it('should clear all client presence keys', async () => { - await service.connect(); - - await service.setClientPresence({ - clientId: '00000000-0000-4000-8000-000000000001', - address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', - clientName: 'worker', - sessionId: '57a6f098', - userAgent: 'bitaxe', - startTime: '2026-06-07T12:00:00.000Z', - lastSeen: '2026-06-07T12:01:00.000Z', - hashRate: 100, - bestDifficulty: 200, - }); - - await service.clearClientPresence(); - - expect(await service.getAllClientPresence()).toHaveLength(0); - expect([...store.keys()].filter(key => key.startsWith('client-presence'))).toHaveLength(0); - expect([...sets.keys()].filter(key => key.startsWith('client-presence'))).toHaveLength(0); - }); - - it('should not return address presence that is no longer in the active global set', async () => { - await service.connect(); - - await service.setClientPresence({ - clientId: '00000000-0000-4000-8000-000000000001', - address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', - clientName: 'worker', - sessionId: '57a6f098', - userAgent: 'bitaxe', - startTime: '2026-06-07T12:00:00.000Z', - lastSeen: '2026-06-07T12:01:00.000Z', - hashRate: 100, - bestDifficulty: 200, - }); - sets.get('client-presence:all')?.delete('00000000-0000-4000-8000-000000000001'); - - expect(await service.getClientPresenceByAddress('tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4')) - .toHaveLength(0); - expect(sets.get('client-presence:address:tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4')) - ?.not.toContain('00000000-0000-4000-8000-000000000001'); - }); }); const store = new Map(); diff --git a/src/services/redis-messaging.service.ts b/src/services/redis-messaging.service.ts index 817e615..0a859cb 100644 --- a/src/services/redis-messaging.service.ts +++ b/src/services/redis-messaging.service.ts @@ -4,45 +4,22 @@ import { createClient, RedisClientType } from 'redis'; import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate'; import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo'; -import { normalizePayoutMode, PayoutMode } from '../types/payout-mode'; const MINING_INFO_CHANNEL = 'mining-info.updated'; const MINING_INFO_KEY = 'mining-info:latest'; const BLOCK_TEMPLATE_LATEST_KEY = 'block-template:latest'; -const CLIENT_PRESENCE_TTL_SECONDS = 180; const blockTemplateKey = (height: number) => `block-template:${height}`; -const CLIENT_PRESENCE_ALL_KEY = 'client-presence:all'; -const clientPresenceKey = (clientId: string) => `client-presence:${clientId}`; -const clientPresenceAddressKey = (address: string) => `client-presence:address:${address}`; const jsonCacheKey = (key: string) => `json-cache:${key}`; -export interface ClientPresence { - clientId: string; - address: string; - clientName: string; - sessionId: string; - payoutMode?: PayoutMode; - userAgent?: string | null; - startTime: string; - lastSeen: string; - hashRate: number; - bestDifficulty: number; -} - @Injectable() export class RedisMessagingService implements OnModuleInit, OnModuleDestroy { private publisher: RedisClientType; private subscriber: RedisClientType; private connected = false; - private readonly clientPresenceTtlSeconds: number; - private readonly clientPresenceEnabled: boolean; constructor( private readonly configService: ConfigService, - ) { - this.clientPresenceTtlSeconds = this.readPositiveInt('CLIENT_PRESENCE_TTL_SECONDS', CLIENT_PRESENCE_TTL_SECONDS); - this.clientPresenceEnabled = this.configService.get('CLIENT_PRESENCE_ENABLED') !== 'false'; - } + ) { } public async onModuleInit() { await this.connect().catch(error => { @@ -140,108 +117,6 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy { return value == null ? null : JSON.parse(value as string); } - public async setClientPresence(presence: ClientPresence): Promise { - if (!this.clientPresenceEnabled) { - return; - } - if (!await this.ensureConnected()) { - return; - } - - const serialized = JSON.stringify({ - ...presence, - payoutMode: normalizePayoutMode(presence.payoutMode), - userAgent: presence.userAgent ?? null, - hashRate: Number.isFinite(Number(presence.hashRate)) ? Number(presence.hashRate) : 0, - bestDifficulty: Number.isFinite(Number(presence.bestDifficulty)) ? Number(presence.bestDifficulty) : 0, - }); - - await Promise.all([ - this.publisher.setEx(clientPresenceKey(presence.clientId), this.clientPresenceTtlSeconds, serialized), - this.publisher.sAdd(CLIENT_PRESENCE_ALL_KEY, presence.clientId), - this.publisher.sAdd(clientPresenceAddressKey(presence.address), presence.clientId), - ]); - } - - public async removeClientPresence(clientId: string, address?: string): Promise { - if (!this.clientPresenceEnabled) { - return; - } - if (!await this.ensureConnected()) { - return; - } - - let resolvedAddress = address; - if (resolvedAddress == null) { - const presence = await this.getClientPresence(clientId); - resolvedAddress = presence?.address; - } - - const removals: Promise[] = [ - this.publisher.del(clientPresenceKey(clientId)), - this.publisher.sRem(CLIENT_PRESENCE_ALL_KEY, clientId), - ]; - if (resolvedAddress != null) { - removals.push(this.publisher.sRem(clientPresenceAddressKey(resolvedAddress), clientId)); - } - - await Promise.all(removals); - } - - public async getClientPresence(clientId: string): Promise { - if (!this.clientPresenceEnabled) { - return null; - } - if (!await this.ensureConnected()) { - return null; - } - - const value = await this.publisher.get(clientPresenceKey(clientId)); - return this.parseClientPresence(value); - } - - public async getClientPresenceByAddress(address: string): Promise { - if (!this.clientPresenceEnabled) { - return []; - } - if (!await this.ensureConnected()) { - return []; - } - - return this.getPresenceFromSet(clientPresenceAddressKey(address)); - } - - public async getAllClientPresence(): Promise { - if (!this.clientPresenceEnabled) { - return []; - } - if (!await this.ensureConnected()) { - return []; - } - - return this.getPresenceFromSet(CLIENT_PRESENCE_ALL_KEY); - } - - public async clearClientPresence(): Promise { - if (!this.clientPresenceEnabled) { - return; - } - if (!await this.ensureConnected()) { - return; - } - - const batch: string[] = []; - for await (const keyOrKeys of (this.publisher as any).scanIterator({ MATCH: 'client-presence*', COUNT: 1000 })) { - const keys = Array.isArray(keyOrKeys) ? keyOrKeys : [keyOrKeys]; - batch.push(...keys.map(key => key as string)); - if (batch.length >= 500) { - await this.deleteKeys(batch.splice(0)); - } - } - - await this.deleteKeys(batch); - } - public async getJsonCache(key: string): Promise { if (!await this.ensureConnected()) { return null; @@ -285,86 +160,4 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy { return true; } - private async getPresenceFromSet(setKey: string): Promise { - const clientIds = await this.publisher.sMembers(setKey); - if (clientIds.length === 0) { - return []; - } - const activeClientIds = setKey === CLIENT_PRESENCE_ALL_KEY - ? null - : new Set(await this.publisher.sMembers(CLIENT_PRESENCE_ALL_KEY)); - - const presences: ClientPresence[] = []; - const staleClientIds: string[] = []; - for (let i = 0; i < clientIds.length; i += 1000) { - const chunk = clientIds.slice(i, i + 1000); - const values = await this.publisher.mGet(chunk.map(clientPresenceKey)); - values.forEach((value, index) => { - const clientId = chunk[index]; - if (activeClientIds != null && !activeClientIds.has(clientId)) { - staleClientIds.push(clientId); - return; - } - const presence = this.parseClientPresence(value); - if (presence == null) { - staleClientIds.push(clientId); - return; - } - presences.push(presence); - }); - } - - for (let i = 0; i < staleClientIds.length; i += 1000) { - const staleChunk = staleClientIds.slice(i, i + 1000); - if (staleChunk.length === 0) { - continue; - } - await this.publisher.sRem(setKey, staleChunk); - if (setKey !== CLIENT_PRESENCE_ALL_KEY) { - await this.publisher.sRem(CLIENT_PRESENCE_ALL_KEY, staleChunk); - } - } - - return presences; - } - - private parseClientPresence(value: unknown): ClientPresence | null { - if (value == null) { - return null; - } - - try { - const parsed = JSON.parse(value as string); - if (parsed?.clientId == null || parsed?.address == null) { - return null; - } - return { - clientId: parsed.clientId, - address: parsed.address, - clientName: parsed.clientName ?? 'default', - sessionId: parsed.sessionId ?? parsed.clientId, - payoutMode: normalizePayoutMode(parsed.payoutMode), - userAgent: parsed.userAgent ?? null, - startTime: parsed.startTime, - lastSeen: parsed.lastSeen, - hashRate: Number(parsed.hashRate ?? 0), - bestDifficulty: Number(parsed.bestDifficulty ?? 0), - }; - } catch (error) { - console.error(`Invalid Redis client presence: ${error.message}`); - return null; - } - } - - private readPositiveInt(name: string, defaultValue: number): number { - const value = Number(this.configService.get(name) ?? process.env[name]); - return Number.isInteger(value) && value > 0 ? value : defaultValue; - } - - private async deleteKeys(keys: string[]): Promise { - if (keys.length === 0) { - return; - } - await (this.publisher as any).del(...keys); - } } diff --git a/src/services/stratum-v1.service.spec.ts b/src/services/stratum-v1.service.spec.ts index 725716d..93f87ce 100644 --- a/src/services/stratum-v1.service.spec.ts +++ b/src/services/stratum-v1.service.spec.ts @@ -31,9 +31,7 @@ describe('StratumV1Service', () => { ensureInitialized: jest.fn().mockResolvedValue(undefined), createClient: jest.fn() }; - redisMessagingService = { - clearClientPresence: jest.fn().mockResolvedValue(undefined) - }; + redisMessagingService = {}; service = new StratumV1Service( {} as any, clientService, @@ -75,7 +73,6 @@ describe('StratumV1Service', () => { jest.runOnlyPendingTimers(); expect(clientService.deleteAll).toHaveBeenCalled(); - expect(redisMessagingService.clearClientPresence).toHaveBeenCalled(); expect(userAgentReportService.refreshReport).toHaveBeenCalled(); expect(startSocketServerSpy).not.toHaveBeenCalled(); expect(startSecureSocketServerSpy).not.toHaveBeenCalled(); diff --git a/src/services/stratum-v1.service.ts b/src/services/stratum-v1.service.ts index 706fa0f..0e6a67f 100644 --- a/src/services/stratum-v1.service.ts +++ b/src/services/stratum-v1.service.ts @@ -81,7 +81,6 @@ export class StratumV1Service implements OnModuleInit { if (process.env.MASTER == 'true') { await this.clientService.deleteAll(); - await this.redisMessagingService?.clearClientPresence(); await this.userAgentReportService.refreshReport(); console.log('Master process skipping Stratum socket listeners'); return; diff --git a/test/timescale-redis.integration-spec.ts b/test/timescale-redis.integration-spec.ts index 6371966..0db0d6d 100644 --- a/test/timescale-redis.integration-spec.ts +++ b/test/timescale-redis.integration-spec.ts @@ -55,7 +55,6 @@ describe('TimescaleDB and Redis integration', () => { await dataSource.query(`DELETE FROM blocks_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'); }); @@ -613,7 +612,7 @@ describe('TimescaleDB and Redis integration', () => { bestDifficulty: 11, hashRate: 0, }); - await repository.save({ + const staleClient = await repository.save({ address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', clientName: 'stale-live-worker', sessionId: '92b2c3d4', @@ -622,18 +621,7 @@ describe('TimescaleDB and Redis integration', () => { bestDifficulty: 22, hashRate: 200, }); - - await redisMessagingService.setClientPresence({ - clientId: activeClient.id, - address: activeClient.address, - clientName: activeClient.clientName, - sessionId: activeClient.sessionId, - userAgent, - startTime: activeClient.startTime.toISOString(), - lastSeen: new Date().toISOString(), - bestDifficulty: Number(activeClient.bestDifficulty), - hashRate: 0, - }); + await repository.softDelete(staleClient.id); const rows = await reportService.getReport(); expect(rows).toEqual(expect.arrayContaining([{