diff --git a/dockerfile b/dockerfile index 18ebaae..3fc7c12 100644 --- a/dockerfile +++ b/dockerfile @@ -3,6 +3,7 @@ FROM node:22-bookworm-slim WORKDIR /app ENV NODE_ENV=production +ENV TZ=Asia/Shanghai ENV APP_PORT=48085 ENV LOG_LEVEL=info ENV LOG_APP_NAME=kt-template-online-api diff --git a/k8s/prod/api.yaml b/k8s/prod/api.yaml index 46c7a57..39f148b 100644 --- a/k8s/prod/api.yaml +++ b/k8s/prod/api.yaml @@ -31,6 +31,8 @@ spec: env: - name: NODE_ENV value: production + - name: TZ + value: Asia/Shanghai - name: QQBOT_NAPCAT_SSH_KEY_PATH value: /app/secrets/napcat-ssh/id_rsa # Jenkins 每次发布会从 Agent 私有 .env.production 重建这个 Secret。 diff --git a/src/admin/system-log/system-log.service.ts b/src/admin/system-log/system-log.service.ts index c46c9bb..3989020 100644 --- a/src/admin/system-log/system-log.service.ts +++ b/src/admin/system-log/system-log.service.ts @@ -30,6 +30,19 @@ type LokiQueryRangeResponse = { status?: string; }; +type LokiMetricResult = { + metric?: Record; + value?: [number | string, string]; +}; + +type LokiQueryResponse = { + data?: { + result?: LokiMetricResult[]; + resultType?: string; + }; + status?: string; +}; + const DEFAULT_LEVELS = ['debug', 'info', 'warning', 'error', 'critical']; const PINO_LEVEL_MAP: Record = { '10': 'debug', @@ -89,17 +102,20 @@ export class SystemLogService { 1, 20, ); + const skip = (pageNo - 1) * pageSize; const requestLimit = Math.min( - this.toolsService.toPositiveNumber(query.limit, pageNo * pageSize), + this.toolsService.toPositiveNumber(query.limit, skip + pageSize), this.getNumberConfig('LOKI_QUERY_MAX_LIMIT', 1000), ); - const logs = await this.queryLogs(query, Math.max(requestLimit, pageSize)); + const [logs, total] = await Promise.all([ + this.queryLogs(query, Math.max(requestLimit, pageSize)), + this.queryLogCount(query), + ]); const filteredLogs = logs.filter((item) => this.matchesQuery(item, query)); - const startIndex = (pageNo - 1) * pageSize; return { - items: filteredLogs.slice(startIndex, startIndex + pageSize), - total: filteredLogs.length, + items: filteredLogs.slice(skip, skip + pageSize), + total, }; } @@ -110,18 +126,10 @@ export class SystemLogService { return DEFAULT_LEVELS.map((level) => ({ count: 0, level })); } - const logs = await this.queryLogs( - { - ...query, - level: undefined, - }, - this.toolsService.toPositiveNumber(query.limit, 1000), - ); - const filteredLogs = logs.filter((item) => this.matchesQuery(item, query)); const countMap = new Map(DEFAULT_LEVELS.map((level) => [level, 0])); - filteredLogs.forEach((item) => { - const level = this.normalizeLevel(item.level) || 'info'; - countMap.set(level, (countMap.get(level) || 0) + 1); + const counts = await this.queryLogSummary(query); + counts.forEach(({ count, level }) => { + countMap.set(level, count); }); return DEFAULT_LEVELS.map((level) => ({ @@ -158,6 +166,51 @@ export class SystemLogService { return this.flattenLogs(response.data?.result || []); } + private async queryLogCount(query: SystemLogQueryDto) { + const response = await this.queryInstant( + this.buildCountLogQL(query), + query, + ); + const value = response.data?.result?.[0]?.value?.[1]; + return this.toOptionalNumber(value) || 0; + } + + private async queryLogSummary(query: SystemLogQueryDto) { + const response = await this.queryInstant( + this.buildSummaryLogQL(query), + query, + ); + + return (response.data?.result || []) + .map((item) => ({ + count: this.toOptionalNumber(item.value?.[1]) || 0, + level: this.normalizeLevel(item.metric?.level) || 'info', + })) + .filter((item) => DEFAULT_LEVELS.includes(item.level)); + } + + private async queryInstant(logql: string, query: SystemLogQueryDto) { + const url = new URL(this.getInstantQueryEndpoint(), this.host); + const { end } = this.getTimeRange(query); + url.searchParams.set('query', logql); + url.searchParams.set('time', `${Math.floor(end.getTime() / 1000)}`); + + let response: LokiQueryResponse; + try { + response = await this.requestJson(url); + } catch (error) { + throwVbenError( + this.toolsService.getErrorMessage(error, 'Loki 查询失败'), + HttpStatus.BAD_GATEWAY, + ); + } + if (response.status && response.status !== 'success') { + throwVbenError('Loki 查询失败', HttpStatus.BAD_GATEWAY, response.status); + } + + return response; + } + private buildLogQL(query: SystemLogQueryDto) { const selector = this.withLevelSelector( this.getBaseSelector(), @@ -176,6 +229,14 @@ export class SystemLogService { return [selector, ...lineFilters].join(' '); } + private buildCountLogQL(query: SystemLogQueryDto) { + return `sum(count_over_time(${this.buildLogQL(query)}[${this.getLogqlRange(query)}]))`; + } + + private buildSummaryLogQL(query: SystemLogQueryDto) { + return `sum by (level)(count_over_time(${this.buildLogQL(query)}[${this.getLogqlRange(query)}]))`; + } + private withLevelSelector(selector: string, level?: string) { if (!level || selector.includes('level=')) return selector; return selector.replace(/}\s*$/, `,level="${this.escapeLabelValue(level)}"}`); @@ -305,6 +366,26 @@ export class SystemLogService { return { end, start }; } + private getLogqlRange(query: SystemLogQueryDto) { + const { end, start } = this.getTimeRange(query); + const seconds = Math.max( + 1, + Math.ceil((end.getTime() - start.getTime()) / 1000), + ); + + return `${seconds}s`; + } + + private getInstantQueryEndpoint() { + const endpoint = this.getConfig('LOKI_QUERY_INSTANT_ENDPOINT'); + if (endpoint) return endpoint; + + return this.getConfig( + 'LOKI_QUERY_ENDPOINT', + '/loki/api/v1/query_range', + ).replace(/query_range$/, 'query'); + } + private requestJson(url: URL) { return new Promise((resolve, reject) => { const client = url.protocol === 'http:' ? http : https; diff --git a/src/common/common.module.ts b/src/common/common.module.ts index 64afa3a..f9b1c9a 100644 --- a/src/common/common.module.ts +++ b/src/common/common.module.ts @@ -1,9 +1,10 @@ import { Global, Module } from '@nestjs/common'; +import { LokiLogPublisherService } from './logger/loki-log-publisher.service'; import { ToolsService } from './services/tool.service'; @Global() @Module({ - exports: [ToolsService], - providers: [ToolsService], + exports: [LokiLogPublisherService, ToolsService], + providers: [LokiLogPublisherService, ToolsService], }) export class CommonModule {} diff --git a/src/common/index.ts b/src/common/index.ts index ca2d381..f13bec2 100644 --- a/src/common/index.ts +++ b/src/common/index.ts @@ -6,6 +6,7 @@ export * from './decorators/public.decorator'; export * from './filters/api-exception.filter'; export * from './interceptors/api-request-log.interceptor'; export * from './interceptors/save-body.interceptor'; +export * from './logger/loki-log-publisher.service'; export * from './logger/pino-logger.config'; export * from './response/vben-response'; export * from './services/tool.service'; diff --git a/src/common/interceptors/api-request-log.interceptor.ts b/src/common/interceptors/api-request-log.interceptor.ts index d49c155..e4f54e1 100644 --- a/src/common/interceptors/api-request-log.interceptor.ts +++ b/src/common/interceptors/api-request-log.interceptor.ts @@ -9,6 +9,7 @@ import { import type { Request, Response } from 'express'; import { PinoLogger } from 'nestjs-pino'; import { catchError, Observable, tap, throwError } from 'rxjs'; +import { LokiLogPublisherService } from '../logger/loki-log-publisher.service'; import { ToolsService } from '../services/tool.service'; type RequestWithId = Request & { @@ -19,6 +20,7 @@ type RequestWithId = Request & { export class ApiRequestLogInterceptor implements NestInterceptor { constructor( private readonly logger: PinoLogger, + private readonly lokiLogPublisherService: LokiLogPublisherService, private readonly toolsService: ToolsService, ) { this.logger.setContext(ApiRequestLogInterceptor.name); @@ -85,6 +87,12 @@ export class ApiRequestLogInterceptor implements NestInterceptor { }; if (statusCode >= 500) { + this.publishRequestLog({ + error: params.error, + level: 'error', + message: 'HTTP request failed', + payload, + }); this.logger.error( { ...payload, @@ -96,13 +104,50 @@ export class ApiRequestLogInterceptor implements NestInterceptor { } if (statusCode >= 400) { + this.publishRequestLog({ + level: 'warning', + message: 'HTTP request completed', + payload, + }); this.logger.warn(payload, 'HTTP request completed'); return; } + this.publishRequestLog({ + level: 'info', + message: 'HTTP request completed', + payload, + }); this.logger.info(payload, 'HTTP request completed'); } + private publishRequestLog(params: { + error?: unknown; + level: 'error' | 'info' | 'warning'; + message: string; + payload: Record; + }) { + if (this.shouldSkipLokiPublish(params.payload.path)) return; + + void this.lokiLogPublisherService + .pushHttpRequestLog({ + context: ApiRequestLogInterceptor.name, + error: params.error, + level: params.level, + message: params.message, + payload: params.payload, + }) + .catch(() => undefined); + } + + private shouldSkipLokiPublish(path: unknown) { + const normalizedPath = this.toolsService.normalizeRequestPathValue(path); + return ( + normalizedPath === '/system/logs' || + normalizedPath.startsWith('/system/logs/') + ); + } + private getStatusCode(error: unknown, response: Response) { if (error instanceof HttpException) { return error.getStatus(); diff --git a/src/common/logger/loki-log-publisher.service.ts b/src/common/logger/loki-log-publisher.service.ts new file mode 100644 index 0000000..9a601ae --- /dev/null +++ b/src/common/logger/loki-log-publisher.service.ts @@ -0,0 +1,174 @@ +import { Injectable } from '@nestjs/common'; +import { ConfigService } from '@nestjs/config'; +import * as http from 'node:http'; +import * as https from 'node:https'; +import { URL } from 'node:url'; +import { ToolsService } from '../services/tool.service'; +import { getAppName, getLokiEnvironment } from './pino-logger.config'; + +type LokiLogLevel = 'critical' | 'debug' | 'error' | 'info' | 'warning'; + +type LokiPushLogParams = { + context: string; + error?: unknown; + level: LokiLogLevel; + message: string; + payload: Record; +}; + +const PINO_LEVEL_VALUES: Record = { + critical: 60, + debug: 20, + error: 50, + info: 30, + warning: 40, +}; + +@Injectable() +export class LokiLogPublisherService { + private readonly appName: string; + private readonly environment: string; + private readonly host: string; + + constructor( + private readonly configService: ConfigService, + private readonly toolsService: ToolsService, + ) { + this.appName = getAppName(configService); + this.environment = getLokiEnvironment(configService); + this.host = this.normalizeUrl( + this.getConfig('LOKI_HOST') || this.getConfig('LOKI_URL'), + ); + } + + async pushHttpRequestLog(params: LokiPushLogParams) { + if (!this.isEnabled()) return; + + const timestampMs = Date.now(); + const stream = { + app: this.appName, + context: params.context, + env: this.environment, + level: params.level, + service: 'api', + }; + const line = JSON.stringify({ + level: PINO_LEVEL_VALUES[params.level], + time: timestampMs, + app: this.appName, + env: this.environment, + context: params.context, + ...params.payload, + ...(params.error ? { err: this.serializeError(params.error) } : {}), + msg: params.message, + }); + const body = JSON.stringify({ + streams: [ + { + stream, + values: [[this.toNanoseconds(timestampMs), line]], + }, + ], + }); + + await this.requestPush(body); + } + + private isEnabled() { + return ( + !!this.host && + this.toolsService.normalizeBoolean( + this.configService.get('LOKI_HTTP_REQUEST_PUSH_ENABLED'), + true, + ) + ); + } + + private requestPush(body: string) { + const url = new URL( + this.getConfig('LOKI_PUSH_ENDPOINT', '/loki/api/v1/push'), + this.host, + ); + + return new Promise((resolve, reject) => { + const client = url.protocol === 'http:' ? http : https; + const request = client.request( + url, + { + headers: { + 'Content-Length': Buffer.byteLength(body), + 'Content-Type': 'application/json', + 'User-Agent': 'kt-template-online-api/loki-log-publisher', + ...this.getHeaders(), + }, + method: 'POST', + timeout: this.getNumberConfig('LOKI_PUSH_TIMEOUT_MS', 30000), + }, + (response) => { + response.resume(); + response.on('end', () => { + if ((response.statusCode || 500) >= 400) { + reject(new Error(`Loki 写入失败:${response.statusCode}`)); + return; + } + resolve(); + }); + }, + ); + + request.on('timeout', () => { + request.destroy(new Error('Loki 写入超时')); + }); + request.on('error', reject); + request.end(body); + }); + } + + private getHeaders() { + const headers: Record = {}; + const tenantId = this.getConfig('LOKI_TENANT_ID'); + const username = this.getConfig('LOKI_USERNAME'); + const password = this.getConfig('LOKI_PASSWORD'); + + if (tenantId) headers['X-Scope-OrgID'] = tenantId; + if (username && password) { + headers.Authorization = `Basic ${Buffer.from( + `${username}:${password}`, + ).toString('base64')}`; + } + + return headers; + } + + private serializeError(error: unknown) { + if (error instanceof Error) { + return { + message: error.message, + name: error.name, + stack: error.stack, + }; + } + + return { + message: this.toolsService.getErrorMessage(error), + }; + } + + private toNanoseconds(timestampMs: number) { + return `${BigInt(timestampMs) * 1000000n}`; + } + + private getConfig(key: string, fallback = '') { + const value = this.configService.get(key); + return this.toolsService.toTrimmedString(value || fallback); + } + + private getNumberConfig(key: string, fallback: number) { + const value = Number(this.configService.get(key)); + return Number.isFinite(value) && value > 0 ? value : fallback; + } + + private normalizeUrl(value: string) { + return value.replace(/\/+$/g, ''); + } +} diff --git a/test/admin/system-log/system-log.service.spec.ts b/test/admin/system-log/system-log.service.spec.ts index 16d756f..c06f0f2 100644 --- a/test/admin/system-log/system-log.service.spec.ts +++ b/test/admin/system-log/system-log.service.spec.ts @@ -17,6 +17,94 @@ function createService() { } describe('SystemLogService', () => { + it('uses Loki aggregate count as page total', async () => { + const service = createService(); + const requestJson = jest + .spyOn(service as any, 'requestJson') + .mockImplementation(async (url: URL) => { + if (url.pathname.endsWith('/query')) { + expect(url.searchParams.get('query')).toContain( + 'sum(count_over_time(', + ); + return { + data: { + result: [ + { + value: [1780576200, '12'], + }, + ], + }, + status: 'success', + }; + } + + return { + data: { + result: [ + { + stream: { + context: 'ApiRequestLogInterceptor', + hostname: 'api-pod', + }, + values: [ + [ + '1780576200000000000', + JSON.stringify({ + durationMs: 18, + level: 30, + method: 'GET', + msg: 'HTTP request completed', + path: '/system/logs', + requestId: 'req-1', + statusCode: 200, + }), + ], + ], + }, + ], + }, + status: 'success', + }; + }); + + const result = await service.page({ + pageNo: 1, + pageSize: 10, + requestId: 'req-1', + }); + + expect(result.total).toBe(12); + expect(result.items).toHaveLength(1); + expect(requestJson).toHaveBeenCalledTimes(2); + }); + + it('uses Loki aggregate counts for summary cards', async () => { + const service = createService(); + jest.spyOn(service as any, 'requestJson').mockResolvedValue({ + data: { + result: [ + { + metric: { level: 'info' }, + value: [1780576200, '8'], + }, + { + metric: { level: 'error' }, + value: [1780576200, '2'], + }, + ], + }, + status: 'success', + }); + + await expect(service.summary({ rangeMinutes: 10 })).resolves.toEqual([ + { count: 0, level: 'debug' }, + { count: 8, level: 'info' }, + { count: 0, level: 'warning' }, + { count: 2, level: 'error' }, + { count: 0, level: 'critical' }, + ]); + }); + it('parses structured HTTP request fields from top-level Loki log lines', () => { const service = createService(); const result = (service as any).serializeLog({ diff --git a/test/common/api-request-log.interceptor.spec.ts b/test/common/api-request-log.interceptor.spec.ts index 3bb0045..42e633c 100644 --- a/test/common/api-request-log.interceptor.spec.ts +++ b/test/common/api-request-log.interceptor.spec.ts @@ -10,7 +10,11 @@ import { Test } from '@nestjs/testing'; import { PinoLogger } from 'nestjs-pino'; import * as request from 'supertest'; import { lastValueFrom, of, throwError } from 'rxjs'; -import { ApiRequestLogInterceptor, ToolsService } from '../../src/common'; +import { + ApiRequestLogInterceptor, + LokiLogPublisherService, + ToolsService, +} from '../../src/common'; @Controller('probe') class ProbeController { @@ -38,6 +42,12 @@ function createResponse(statusCode = 200) { }; } +function createLokiLogPublisherMock() { + return { + pushHttpRequestLog: jest.fn().mockResolvedValue(undefined), + }; +} + describe('ApiRequestLogInterceptor', () => { it('writes structured HTTP request fields for successful controller calls', async () => { const logger = { @@ -46,8 +56,10 @@ describe('ApiRequestLogInterceptor', () => { setContext: jest.fn(), warn: jest.fn(), }; + const lokiLogPublisher = createLokiLogPublisherMock(); const interceptor = new ApiRequestLogInterceptor( logger as any, + lokiLogPublisher as any, new ToolsService(), ); const request: Record = { @@ -55,7 +67,7 @@ describe('ApiRequestLogInterceptor', () => { 'x-request-id': 'req-1', }, method: 'GET', - originalUrl: '/system/logs?pageNo=1', + originalUrl: '/status?source=test', }; const response = createResponse(200); @@ -70,12 +82,65 @@ describe('ApiRequestLogInterceptor', () => { expect(logger.info).toHaveBeenCalledWith( expect.objectContaining({ method: 'GET', - path: '/system/logs', + path: '/status', requestId: 'req-1', statusCode: 200, }), 'HTTP request completed', ); + expect(lokiLogPublisher.pushHttpRequestLog).toHaveBeenCalledWith( + expect.objectContaining({ + context: ApiRequestLogInterceptor.name, + level: 'info', + message: 'HTTP request completed', + payload: expect.objectContaining({ + method: 'GET', + path: '/status', + requestId: 'req-1', + statusCode: 200, + }), + }), + ); + }); + + it('does not publish system log query requests back to Loki', async () => { + const logger = { + error: jest.fn(), + info: jest.fn(), + setContext: jest.fn(), + warn: jest.fn(), + }; + const lokiLogPublisher = createLokiLogPublisherMock(); + const interceptor = new ApiRequestLogInterceptor( + logger as any, + lokiLogPublisher as any, + new ToolsService(), + ); + const request: Record = { + headers: { + 'x-request-id': 'req-system-log', + }, + method: 'GET', + originalUrl: '/system/logs?pageNo=1', + }; + const response = createResponse(200); + + await lastValueFrom( + interceptor.intercept(createHttpContext(request, response), { + handle: () => of({ ok: true }), + }), + ); + + expect(logger.info).toHaveBeenCalledWith( + expect.objectContaining({ + method: 'GET', + path: '/system/logs', + requestId: 'req-system-log', + statusCode: 200, + }), + 'HTTP request completed', + ); + expect(lokiLogPublisher.pushHttpRequestLog).not.toHaveBeenCalled(); }); it('writes warning logs with HTTP status for controller errors', async () => { @@ -85,8 +150,10 @@ describe('ApiRequestLogInterceptor', () => { setContext: jest.fn(), warn: jest.fn(), }; + const lokiLogPublisher = createLokiLogPublisherMock(); const interceptor = new ApiRequestLogInterceptor( logger as any, + lokiLogPublisher as any, new ToolsService(), ); const request: Record = { @@ -115,6 +182,17 @@ describe('ApiRequestLogInterceptor', () => { }), 'HTTP request completed', ); + expect(lokiLogPublisher.pushHttpRequestLog).toHaveBeenCalledWith( + expect.objectContaining({ + level: 'warning', + payload: expect.objectContaining({ + method: 'POST', + path: '/dict/save', + requestId: request.id, + statusCode: 404, + }), + }), + ); }); it('captures a real local HTTP request in a Nest application', async () => { @@ -135,6 +213,10 @@ describe('ApiRequestLogInterceptor', () => { provide: PinoLogger, useValue: logger, }, + { + provide: LokiLogPublisherService, + useValue: createLokiLogPublisherMock(), + }, { provide: APP_INTERCEPTOR, useClass: ApiRequestLogInterceptor,