Add API timing diagnostics

This commit is contained in:
Ben
2026-06-08 00:47:11 -04:00
parent df35549f27
commit 7b7dc87242
7 changed files with 149 additions and 38 deletions
@@ -5,6 +5,7 @@ import { Repository } from 'typeorm';
import { ClientEntity } from '../../client/client.entity'; import { ClientEntity } from '../../client/client.entity';
import { UserAgentReportView } from './user-agent-report.view'; import { UserAgentReportView } from './user-agent-report.view';
import { RedisMessagingService } from '../../../services/redis-messaging.service'; import { RedisMessagingService } from '../../../services/redis-messaging.service';
import { logTiming, timeAsync, timingStart } from '../../../utils/timing.utils';
@Injectable() @Injectable()
export class UserAgentReportService { export class UserAgentReportService {
@@ -22,21 +23,27 @@ export class UserAgentReportService {
} }
public async getReport() { public async getReport() {
const cachedReport = await this.redisMessagingService const start = timingStart();
const cachedReport = await timeAsync('user agent report shared cache read', () => this.redisMessagingService
.getJsonCache<UserAgentReportView[]>(this.liveReportCacheKey) .getJsonCache<UserAgentReportView[]>(this.liveReportCacheKey)
.catch(error => { .catch(error => {
console.error(`Live user-agent report cache read failed: ${error.message}`); console.error(`Live user-agent report cache read failed: ${error.message}`);
return null; return null;
}); }));
if (cachedReport != null) { if (cachedReport != null) {
logTiming('user agent report getReport', start, { cache: 'hit', rows: cachedReport.length });
return cachedReport; return cachedReport;
} }
if (process.env.API_ONLY == 'true') { if (process.env.API_ONLY == 'true') {
return this.userAgentReport.find(); const rows = await timeAsync('user agent report materialized view read', () => this.userAgentReport.find());
logTiming('user agent report getReport', start, { cache: 'miss', source: 'materialized-view', rows: rows.length });
return rows;
} }
return this.refreshLiveReport(); const rows = await this.refreshLiveReport();
logTiming('user agent report getReport', start, { cache: 'miss', source: 'live-presence', rows: rows.length });
return rows;
} }
public async refreshLiveReport() { public async refreshLiveReport() {
@@ -53,7 +60,8 @@ export class UserAgentReportService {
} }
private async buildLiveReport() { private async buildLiveReport() {
const presences = await this.redisMessagingService.getAllClientPresence(); const start = timingStart();
const presences = await timeAsync('user agent report all presence load', () => this.redisMessagingService.getAllClientPresence());
const rows = new Map<string, { const rows = new Map<string, {
userAgent: string; userAgent: string;
count: number; count: number;
@@ -95,6 +103,7 @@ export class UserAgentReportService {
console.error(`Live user-agent report cache write failed: ${error.message}`); console.error(`Live user-agent report cache write failed: ${error.message}`);
}); });
logTiming('user agent report buildLiveReport', start, { presences: presences.length, rows: report.length });
return report; return report;
} }
@@ -1,6 +1,7 @@
import { Injectable } from '@nestjs/common'; import { Injectable } from '@nestjs/common';
import { InjectDataSource } from '@nestjs/typeorm'; import { InjectDataSource } from '@nestjs/typeorm';
import { DataSource } from 'typeorm'; import { DataSource } from 'typeorm';
import { timeAsync } from '../../utils/timing.utils';
const HASHES_PER_DIFFICULTY = 4294967296; const HASHES_PER_DIFFICULTY = 4294967296;
const CHART_BUCKET_SECONDS = 600; const CHART_BUCKET_SECONDS = 600;
@@ -31,14 +32,14 @@ export class ClientStatisticsService {
} }
public async getHashRateForGroup(address: string, clientName: string) { public async getHashRateForGroup(address: string, clientName: string) {
const result = await this.dataSource.query(` const result = await timeAsync('client statistics group hashrate query', () => this.dataSource.query(`
SELECT SELECT
COALESCE((SUM("creditedDifficulty") * ${HASHES_PER_DIFFICULTY}) / ${CHART_BUCKET_SECONDS}, 0) AS "hashRate" COALESCE((SUM("creditedDifficulty") * ${HASHES_PER_DIFFICULTY}) / ${CHART_BUCKET_SECONDS}, 0) AS "hashRate"
FROM "accepted_share_entity" FROM "accepted_share_entity"
WHERE "address" = $1 WHERE "address" = $1
AND "clientName" = $2 AND "clientName" = $2
AND "acceptedAt" > NOW() - INTERVAL '1 hour' AND "acceptedAt" > NOW() - INTERVAL '1 hour'
`, [address, clientName]); `, [address, clientName]), { address, clientName });
return parseFloat(result[0]?.hashRate ?? '0'); return parseFloat(result[0]?.hashRate ?? '0');
} }
@@ -90,7 +91,12 @@ export class ClientStatisticsService {
ORDER BY "label" ORDER BY "label"
`; `;
const result = await this.dataSource.query(query, params); const result = await timeAsync('client statistics chart query', () => this.dataSource.query(query, params), {
filterSql,
params: params.length,
limit,
windowSql,
});
return result.map(res => { return result.map(res => {
return { return {
@@ -4,6 +4,7 @@ import { Repository } from 'typeorm';
import { AcceptedShareEntity } from '../accepted-share/accepted-share.entity'; import { AcceptedShareEntity } from '../accepted-share/accepted-share.entity';
import { RedisMessagingService } from '../../services/redis-messaging.service'; import { RedisMessagingService } from '../../services/redis-messaging.service';
import { timeAsync } from '../../utils/timing.utils';
export interface AcceptedShareRecord { export interface AcceptedShareRecord {
protocol: 'sv1' | 'sv2'; protocol: 'sv1' | 'sv2';
@@ -195,7 +196,7 @@ export class ShareAccountingService implements OnModuleDestroy {
return summaries; return summaries;
} }
const rows = await this.acceptedShareRepository.query(` const rows = await timeAsync('share accounting session summaries query', () => this.acceptedShareRepository.query(`
SELECT SELECT
"clientId", "clientId",
MAX("acceptedAt") AS "latestShareAt", MAX("acceptedAt") AS "latestShareAt",
@@ -204,7 +205,7 @@ export class ShareAccountingService implements OnModuleDestroy {
FROM "accepted_share_entity" FROM "accepted_share_entity"
WHERE "clientId" = ANY($1::uuid[]) WHERE "clientId" = ANY($1::uuid[])
GROUP BY "clientId" GROUP BY "clientId"
`, [uniqueClientIds]); `, [uniqueClientIds]), { clientIds: uniqueClientIds.length });
rows.forEach(row => { rows.forEach(row => {
summaries.set(row.clientId, { summaries.set(row.clientId, {
@@ -247,7 +248,7 @@ export class ShareAccountingService implements OnModuleDestroy {
private async loadSummary(filter: AccountingFilter): Promise<ShareAccountingSummary> { private async loadSummary(filter: AccountingFilter): Promise<ShareAccountingSummary> {
const { whereSql, params } = this.buildWhereClause(filter); const { whereSql, params } = this.buildWhereClause(filter);
const [summary] = await this.acceptedShareRepository.query(` const [summary] = await timeAsync('share accounting summary query', () => this.acceptedShareRepository.query(`
SELECT SELECT
COUNT(*)::int AS "totalAcceptedShares", COUNT(*)::int AS "totalAcceptedShares",
COALESCE(SUM("creditedDifficulty"), 0)::float AS "totalCreditedDifficulty", COALESCE(SUM("creditedDifficulty"), 0)::float AS "totalCreditedDifficulty",
@@ -264,9 +265,9 @@ export class ShareAccountingService implements OnModuleDestroy {
MAX("acceptedAt") AS "latestShareAt" MAX("acceptedAt") AS "latestShareAt"
FROM "accepted_share_entity" FROM "accepted_share_entity"
${whereSql} ${whereSql}
`, params); `, params), { filter, params: params.length });
const protocolRows = await this.acceptedShareRepository.query(` const protocolRows = await timeAsync('share accounting protocol breakdown query', () => this.acceptedShareRepository.query(`
SELECT SELECT
"protocol", "protocol",
COUNT(*)::int AS "acceptedShares", COUNT(*)::int AS "acceptedShares",
@@ -275,7 +276,7 @@ export class ShareAccountingService implements OnModuleDestroy {
${whereSql} ${whereSql}
GROUP BY "protocol" GROUP BY "protocol"
ORDER BY "protocol" ORDER BY "protocol"
`, params); `, params), { filter, params: params.length });
return { return {
totalAcceptedShares: this.toNumber(summary?.totalAcceptedShares), totalAcceptedShares: this.toNumber(summary?.totalAcceptedShares),
+26 -9
View File
@@ -13,6 +13,7 @@ import { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-r
import { StratumV2Service } from './services/stratum-v2.service'; import { StratumV2Service } from './services/stratum-v2.service';
import { ShareAccountingService } from './ORM/share-accounting/share-accounting.service'; import { ShareAccountingService } from './ORM/share-accounting/share-accounting.service';
import { RedisMessagingService } from './services/redis-messaging.service'; import { RedisMessagingService } from './services/redis-messaging.service';
import { logTiming, timeAsync, timingStart } from './utils/timing.utils';
@Controller() @Controller()
export class AppController { export class AppController {
@@ -35,6 +36,7 @@ export class AppController {
@Get('info') @Get('info')
public async info() { public async info() {
const start = timingStart();
const CACHE_KEY = 'SITE_INFO'; const CACHE_KEY = 'SITE_INFO';
@@ -42,16 +44,20 @@ export class AppController {
const cachedResult = await this.getCached(CACHE_KEY, 5 * 60 * 1000); const cachedResult = await this.getCached(CACHE_KEY, 5 * 60 * 1000);
if (cachedResult != null) { if (cachedResult != null) {
logTiming('GET /api/info', start, { cache: 'fresh' });
return cachedResult; return cachedResult;
} }
const staleResult = await this.getCached<SiteInfoResponse>(STALE_CACHE_KEY, 60 * 60 * 1000); const staleResult = await this.getCached<SiteInfoResponse>(STALE_CACHE_KEY, 60 * 60 * 1000);
if (staleResult != null) { if (staleResult != null) {
void this.refreshSiteInfo(staleResult); void this.refreshSiteInfo(staleResult);
logTiming('GET /api/info', start, { cache: 'stale' });
return staleResult; return staleResult;
} }
return this.refreshSiteInfo(null); const response = await this.refreshSiteInfo(null);
logTiming('GET /api/info', start, { cache: 'miss' });
return response;
} }
@@ -76,13 +82,13 @@ export class AppController {
}; };
const [blockData, highScores, poolAuthority, userAgentReport] = await Promise.all([ const [blockData, highScores, poolAuthority, userAgentReport] = await Promise.all([
withInfoTimeout('found blocks', this.blocksService.getFoundBlocks(), staleInfo?.blockData ?? []), withInfoTimeout('found blocks', timeAsync('/api/info found blocks', () => this.blocksService.getFoundBlocks()), staleInfo?.blockData ?? []),
withInfoTimeout('high scores', this.addressSettingsService.getHighScores(), staleInfo?.highScores ?? []), withInfoTimeout('high scores', timeAsync('/api/info high scores', () => this.addressSettingsService.getHighScores()), staleInfo?.highScores ?? []),
withInfoTimeout('SV2 authority', this.stratumV2Service.getPoolAuthorityPublicKey(), { withInfoTimeout('SV2 authority', timeAsync('/api/info SV2 authority', () => this.stratumV2Service.getPoolAuthorityPublicKey()), {
publicKey: staleInfo?.sv2?.poolAuthorityPublicKey ?? '', publicKey: staleInfo?.sv2?.poolAuthorityPublicKey ?? '',
configured: staleInfo?.sv2?.authorityKeyConfigured ?? false configured: staleInfo?.sv2?.authorityKeyConfigured ?? false
}), }),
withInfoTimeout<UserAgentReportView[]>('user agent report', this.userAgentReportService.getReport(), staleInfo?.userAgents ?? []), withInfoTimeout<UserAgentReportView[]>('user agent report', timeAsync('/api/info user agent report', () => this.userAgentReportService.getReport()), staleInfo?.userAgents ?? []),
]); ]);
const other: { const other: {
@@ -133,37 +139,42 @@ export class AppController {
@Get('info/accounting') @Get('info/accounting')
public async infoAccounting() { public async infoAccounting() {
const start = timingStart();
const CACHE_KEY = 'SITE_ACCOUNTING'; const CACHE_KEY = 'SITE_ACCOUNTING';
const cachedResult = await this.getCached(CACHE_KEY, 15 * 1000); const cachedResult = await this.getCached(CACHE_KEY, 15 * 1000);
if (cachedResult != null) { if (cachedResult != null) {
logTiming('GET /api/info/accounting', start, { cache: 'hit' });
return cachedResult; return cachedResult;
} }
const data = await this.shareAccountingService.getPoolSummary(); const data = await timeAsync('/api/info/accounting getPoolSummary', () => this.shareAccountingService.getPoolSummary());
//15 sec //15 sec
await this.setCached(CACHE_KEY, data, 15 * 1000); await this.setCached(CACHE_KEY, data, 15 * 1000);
logTiming('GET /api/info/accounting', start, { cache: 'miss' });
return data; return data;
} }
@Get('pool') @Get('pool')
public async pool() { public async pool() {
const start = timingStart();
const CACHE_KEY = 'POOL_INFO'; const CACHE_KEY = 'POOL_INFO';
const cachedResult = await this.getCached(CACHE_KEY, 15 * 1000); const cachedResult = await this.getCached(CACHE_KEY, 15 * 1000);
if (cachedResult != null) { if (cachedResult != null) {
logTiming('GET /api/pool', start, { cache: 'hit' });
return cachedResult; return cachedResult;
} }
const userAgents = await this.userAgentReportService.getReport(); const userAgents = await timeAsync('/api/pool user agent report', () => this.userAgentReportService.getReport());
const totalHashRate = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.totalHashRate), 0); const totalHashRate = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.totalHashRate), 0);
const totalMiners = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.count), 0); const totalMiners = userAgents.reduce((acc, userAgent) => acc + parseFloat(userAgent.count), 0);
const blockHeight = this.bitcoinRpcService.miningInfo.blocks; const blockHeight = this.bitcoinRpcService.miningInfo.blocks;
const blocksFound = await this.blocksService.getFoundBlocks(); const blocksFound = await timeAsync('/api/pool found blocks', () => this.blocksService.getFoundBlocks());
const data = { const data = {
totalHashRate, totalHashRate,
@@ -176,30 +187,36 @@ export class AppController {
// Keep online miner counts responsive after reconnect cleanup. // Keep online miner counts responsive after reconnect cleanup.
await this.setCached(CACHE_KEY, data, 15 * 1000); await this.setCached(CACHE_KEY, data, 15 * 1000);
logTiming('GET /api/pool', start, { cache: 'miss', userAgentCount: userAgents.length });
return data; return data;
} }
@Get('network') @Get('network')
public async network() { public async network() {
const start = timingStart();
logTiming('GET /api/network', start);
return this.bitcoinRpcService.miningInfo ?? {}; return this.bitcoinRpcService.miningInfo ?? {};
} }
@Get('info/chart') @Get('info/chart')
public async infoChart() { public async infoChart() {
const start = timingStart();
const CACHE_KEY = 'SITE_HASHRATE_GRAPH'; const CACHE_KEY = 'SITE_HASHRATE_GRAPH';
const cachedResult = await this.getCached(CACHE_KEY, 10 * 60 * 1000); const cachedResult = await this.getCached(CACHE_KEY, 10 * 60 * 1000);
if (cachedResult != null) { if (cachedResult != null) {
logTiming('GET /api/info/chart', start, { cache: 'hit' });
return cachedResult; return cachedResult;
} }
const chartData = await this.clientStatisticsService.getChartDataForSite(); const chartData = await timeAsync('/api/info/chart query', () => this.clientStatisticsService.getChartDataForSite());
//10 min //10 min
await this.setCached(CACHE_KEY, chartData, 10 * 60 * 1000); await this.setCached(CACHE_KEY, chartData, 10 * 60 * 1000);
logTiming('GET /api/info/chart', start, { cache: 'miss', points: chartData.length });
return chartData; return chartData;
+33 -15
View File
@@ -5,6 +5,7 @@ import { ClientStatisticsService } from '../../ORM/client-statistics/client-stat
import { ClientService } from '../../ORM/client/client.service'; import { ClientService } from '../../ORM/client/client.service';
import { ShareAccountingService } from '../../ORM/share-accounting/share-accounting.service'; import { ShareAccountingService } from '../../ORM/share-accounting/share-accounting.service';
import { RedisMessagingService } from '../../services/redis-messaging.service'; import { RedisMessagingService } from '../../services/redis-messaging.service';
import { logTiming, timeAsync, timingStart } from '../../utils/timing.utils';
@Controller('client') @Controller('client')
@@ -21,16 +22,18 @@ export class ClientController {
@Get(':address') @Get(':address')
async getClientInfo(@Param('address') address: string) { async getClientInfo(@Param('address') address: string) {
const start = timingStart();
const workers = await this.redisMessagingService.getClientPresenceByAddress(address); const workers = await timeAsync('/api/client/:address presence', () => this.redisMessagingService.getClientPresenceByAddress(address), { address });
const sessionSummaries = await this.shareAccountingService.getSessionSummaries(workers.map(worker => worker.clientId)); const sessionSummaries = await timeAsync('/api/client/:address session summaries', () => this.shareAccountingService.getSessionSummaries(workers.map(worker => worker.clientId)), { address, workers: workers.length });
const addressSettings = await this.addressSettingsService.getSettings(address, false); const addressSettings = 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 });
return { const response = {
bestDifficulty: addressSettings?.bestDifficulty, bestDifficulty: addressSettings?.bestDifficulty,
workersCount: workers.length, workersCount: workers.length,
accounting: await this.shareAccountingService.getAddressSummary(address), accounting,
workers: await Promise.all( workers: await Promise.all(
workers.map(async (worker) => { workers.map(async (worker) => {
const sessionSummary = sessionSummaries.get(worker.clientId); const sessionSummary = sessionSummaries.get(worker.clientId);
@@ -49,18 +52,24 @@ export class ClientController {
}) })
) )
} }
logTiming('GET /api/client/:address', start, { address, workers: workers.length });
return response;
} }
@Get(':address/chart') @Get(':address/chart')
async getClientInfoChart(@Param('address') address: string) { async getClientInfoChart(@Param('address') address: string) {
const chartData = await this.clientStatisticsService.getChartDataForAddress(address); const start = timingStart();
const chartData = await timeAsync('/api/client/:address/chart query', () => this.clientStatisticsService.getChartDataForAddress(address), { address });
logTiming('GET /api/client/:address/chart', start, { address, points: chartData.length });
return chartData; return chartData;
} }
@Get(':address/:workerName') @Get(':address/:workerName')
async getWorkerGroupInfo(@Param('address') address: string, @Param('workerName') workerName: string) { async getWorkerGroupInfo(@Param('address') address: string, @Param('workerName') workerName: string) {
const start = timingStart();
const workers = (await this.redisMessagingService.getClientPresenceByAddress(address)) const addressWorkers = await timeAsync('/api/client/:address/:workerName address presence', () => this.redisMessagingService.getClientPresenceByAddress(address), { address, workerName });
const workers = addressWorkers
.filter(worker => worker.clientName === workerName); .filter(worker => worker.clientName === workerName);
const bestDifficulty = workers.reduce((pre, cur, idx, arr) => { const bestDifficulty = workers.reduce((pre, cur, idx, arr) => {
@@ -70,24 +79,29 @@ export class ClientController {
return pre; return pre;
}, 0); }, 0);
const chartData = await this.clientStatisticsService.getChartDataForGroup(address, workerName); const chartData = await timeAsync('/api/client/:address/:workerName chart', () => this.clientStatisticsService.getChartDataForGroup(address, workerName), { address, workerName });
return { const accounting = await timeAsync('/api/client/:address/:workerName accounting', () => this.shareAccountingService.getWorkerGroupSummary(address, workerName), { address, workerName });
const response = {
name: workerName, name: workerName,
bestDifficulty: Math.floor(bestDifficulty), bestDifficulty: Math.floor(bestDifficulty),
accounting: await this.shareAccountingService.getWorkerGroupSummary(address, workerName), accounting,
chartData: chartData, 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') @Get(':address/:workerName/:sessionId')
async getWorkerInfo(@Param('address') address: string, @Param('workerName') workerName: string, @Param('sessionId') sessionId: string) { async getWorkerInfo(@Param('address') address: string, @Param('workerName') workerName: string, @Param('sessionId') sessionId: string) {
const start = timingStart();
const presenceWorker = (await this.redisMessagingService.getClientPresenceByAddress(address)) const addressWorkers = await timeAsync('/api/client/:address/:workerName/:sessionId address presence', () => this.redisMessagingService.getClientPresenceByAddress(address), { address, workerName, sessionId });
const presenceWorker = addressWorkers
.find(worker => worker.clientName === workerName && worker.sessionId === sessionId); .find(worker => worker.clientName === workerName && worker.sessionId === sessionId);
const worker = presenceWorker == null const worker = presenceWorker == null
? await this.clientService.getBySessionId(address, workerName, sessionId) ? await timeAsync('/api/client/:address/:workerName/:sessionId DB fallback', () => this.clientService.getBySessionId(address, workerName, sessionId), { address, workerName, sessionId })
: { : {
id: presenceWorker.clientId, id: presenceWorker.clientId,
sessionId: presenceWorker.sessionId, sessionId: presenceWorker.sessionId,
@@ -96,17 +110,21 @@ export class ClientController {
startTime: presenceWorker.startTime, startTime: presenceWorker.startTime,
}; };
if (worker == null) { if (worker == null) {
logTiming('GET /api/client/:address/:workerName/:sessionId', start, { address, workerName, sessionId, found: false, addressWorkers: addressWorkers.length });
return new NotFoundException(); return new NotFoundException();
} }
const chartData = await this.clientStatisticsService.getChartDataForSession(worker.id); const chartData = await timeAsync('/api/client/:address/:workerName/:sessionId chart', () => this.clientStatisticsService.getChartDataForSession(worker.id), { address, workerName, sessionId, clientId: worker.id });
const accounting = await timeAsync('/api/client/:address/:workerName/:sessionId accounting', () => this.shareAccountingService.getSessionSummary(worker.id), { address, workerName, sessionId, clientId: worker.id });
return { const response = {
sessionId: worker.sessionId, sessionId: worker.sessionId,
name: worker.clientName, name: worker.clientName,
bestDifficulty: Math.floor(worker.bestDifficulty), bestDifficulty: Math.floor(worker.bestDifficulty),
accounting: await this.shareAccountingService.getSessionSummary(worker.id), accounting,
chartData: chartData, chartData: chartData,
startTime: worker.startTime startTime: worker.startTime
} }
logTiming('GET /api/client/:address/:workerName/:sessionId', start, { address, workerName, sessionId, found: true, addressWorkers: addressWorkers.length, points: chartData.length });
return response;
} }
} }
+9
View File
@@ -4,6 +4,7 @@ import { createClient, RedisClientType } from 'redis';
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate'; import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo'; import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo';
import { logTiming, timingStart } from '../utils/timing.utils';
const MINING_INFO_CHANNEL = 'mining-info.updated'; const MINING_INFO_CHANNEL = 'mining-info.updated';
const MINING_INFO_KEY = 'mining-info:latest'; const MINING_INFO_KEY = 'mining-info:latest';
@@ -263,8 +264,10 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
} }
private async getPresenceFromSet(setKey: string): Promise<ClientPresence[]> { private async getPresenceFromSet(setKey: string): Promise<ClientPresence[]> {
const start = timingStart();
const clientIds = await this.publisher.sMembers(setKey); const clientIds = await this.publisher.sMembers(setKey);
if (clientIds.length === 0) { if (clientIds.length === 0) {
logTiming('redis presence set load', start, { setKey, clientIds: 0, presences: 0, staleClientIds: 0 });
return []; return [];
} }
@@ -294,6 +297,12 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
} }
} }
logTiming('redis presence set load', start, {
setKey,
clientIds: clientIds.length,
presences: presences.length,
staleClientIds: staleClientIds.length,
});
return presences; return presences;
} }
+51
View File
@@ -0,0 +1,51 @@
import { performance } from 'perf_hooks';
type TimingMetadata = Record<string, unknown> | (() => Record<string, unknown>);
const DEFAULT_TIMING_LOG_MS = 250;
export function timingStart(): number {
return performance.now();
}
export async function timeAsync<T>(
label: string,
work: () => Promise<T>,
metadata?: TimingMetadata,
): Promise<T> {
const start = timingStart();
try {
return await work();
} finally {
logTiming(label, start, metadata);
}
}
export function logTiming(label: string, start: number, metadata?: TimingMetadata): void {
const elapsedMs = performance.now() - start;
const thresholdMs = getTimingThresholdMs();
if (thresholdMs < 0 || elapsedMs < thresholdMs) {
return;
}
const resolvedMetadata = resolveMetadata(metadata);
const metadataText = resolvedMetadata == null ? '' : ` ${JSON.stringify(resolvedMetadata)}`;
console.warn(`[timing] ${label} ${elapsedMs.toFixed(1)}ms${metadataText}`);
}
function getTimingThresholdMs(): number {
const configured = Number(process.env.API_TIMING_LOG_MS);
return Number.isFinite(configured) ? configured : DEFAULT_TIMING_LOG_MS;
}
function resolveMetadata(metadata?: TimingMetadata): Record<string, unknown> | null {
if (metadata == null) {
return null;
}
try {
return typeof metadata === 'function' ? metadata() : metadata;
} catch (error) {
return { metadataError: error instanceof Error ? error.message : String(error) };
}
}