fix: 修复系统日志Loki链路

This commit is contained in:
sunlei 2026-06-04 21:07:13 +08:00
parent b1249517e8
commit f7c0fb2aff
9 changed files with 496 additions and 21 deletions

View File

@ -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

View File

@ -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。

View File

@ -30,6 +30,19 @@ type LokiQueryRangeResponse = {
status?: string;
};
type LokiMetricResult = {
metric?: Record<string, string>;
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<string, string> = {
'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<LokiQueryResponse>(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<T>(url: URL) {
return new Promise<T>((resolve, reject) => {
const client = url.protocol === 'http:' ? http : https;

View File

@ -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 {}

View File

@ -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';

View File

@ -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<string, unknown>;
}) {
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();

View File

@ -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<string, unknown>;
};
const PINO_LEVEL_VALUES: Record<LokiLogLevel, number> = {
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<string>('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<void>((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<string, string> = {};
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<string>(key);
return this.toolsService.toTrimmedString(value || fallback);
}
private getNumberConfig(key: string, fallback: number) {
const value = Number(this.configService.get<string>(key));
return Number.isFinite(value) && value > 0 ? value : fallback;
}
private normalizeUrl(value: string) {
return value.replace(/\/+$/g, '');
}
}

View File

@ -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({

View File

@ -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<string, any> = {
@ -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<string, any> = {
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<string, any> = {
@ -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,