feat: add privacy-safe server diagnostics

This commit is contained in:
inman committed 2026-08-30 16:01:42 +08:00
1 parent 2360506607
commit a3963ad0f6
27 files changed
+1305 -112

No files matched your search

+341 -28
View File
@@ -1,12 +1,14 @@
import { randomUUID } from 'node:crypto';
import { hostname } from 'node:os';
import { pathToFileURL } from 'node:url';
import { fileURLToPath } from 'node:url';
import { resolve } from 'node:path';
import Fastify, { type FastifyReply, type FastifyRequest } from 'fastify';
import Fastify, { LogController, type FastifyReply, type FastifyRequest } from 'fastify';
import cookie from '@fastify/cookie';
import helmet from '@fastify/helmet';
import rateLimit from '@fastify/rate-limit';
import fastifyStatic from '@fastify/static';
import pino, { type DestinationStream } from 'pino';
import { z } from 'zod';
import { loadConfig, type AppConfig } from './config.js';
import { assertDatabaseSchema, closePool, databaseReadiness, databaseReady, getPool } from './db.js';
@@ -31,6 +33,14 @@ import {
type EncodedTaskInputAttachment,
type TaskInputAttachmentInput
} from './input-attachment.js';
import {
diagnosticDurationMs,
diagnosticError,
diagnosticMetadataKeys,
diagnosticRequestPath,
normalizeRequestId,
writeEmergencyDiagnostic
} from './diagnostics.js';
export function aiServiceConnected(databaseIsReady: boolean, probe: Record<string, unknown>): boolean {
return databaseIsReady && probe.configured === true && probe.reachable === true;
@@ -120,8 +130,64 @@ const taskBulkDeleteSchema = z.object({
{ message: '任务编号不能重复。', path: ['task_ids'] }
);
const LOG_REDACTION_PATHS = [
'req.headers.authorization',
'req.headers.cookie',
'req.headers["x-csrf-token"]',
'request.headers.authorization',
'request.headers.cookie',
'authorization',
'cookie',
'password',
'api_key',
'token',
'access_key',
'*.authorization',
'*.cookie',
'*.password',
'*.api_key',
'*.token',
'*.access_key'
];
export function createControlPlaneLogger(config: AppConfig, destination?: DestinationStream) {
const options = {
level: config.LOG_LEVEL,
base: {
pid: process.pid,
hostname: hostname(),
service: 'ltjt-control-plane',
environment: config.NODE_ENV,
deployment_revision: config.DEPLOYMENT_REVISION
},
redact: {
paths: LOG_REDACTION_PATHS,
censor: '[REDACTED]'
}
};
return destination ? pino(options, destination) : pino(options);
}
function agentBusDiagnosticMetadata(metadata: Record<string, unknown>): Record<string, unknown> {
const rawEvent = String(metadata.agentbus_event || 'event');
const event = /^[a-z0-9_]{1,80}$/u.test(rawEvent) ? rawEvent : 'event';
return {
...metadata,
diagnostic_event: `agentbus.${event}`,
diagnostic_stage: 'agentbus'
};
}
function requestId(request: FastifyRequest): string {
return String(request.headers['x-request-id'] || randomUUID());
return String(request.id);
}
function requestTaskId(request: FastifyRequest): string | undefined {
const params = request.params && typeof request.params === 'object'
? request.params as Record<string, unknown>
: {};
const taskId = String(params.taskId || '').trim();
return taskId ? taskId.slice(0, 200) : undefined;
}
function clientAddress(request: FastifyRequest): string {
@@ -251,15 +317,19 @@ async function loadExternalParser(): Promise<ExternalParser> {
export async function buildServer({
config = loadConfig(),
parser,
startParserLoop = true
startParserLoop = true,
loggerDestination
}: {
config?: AppConfig;
parser?: ExternalParser;
startParserLoop?: boolean;
loggerDestination?: DestinationStream;
} = {}) {
const requestStartedAt = new WeakMap<FastifyRequest, bigint>();
const app = Fastify({
logger: { level: config.LOG_LEVEL, redact: ['req.headers.cookie', 'req.headers.authorization', '*.password', '*.api_key', '*.token'] },
requestIdHeader: 'x-request-id',
loggerInstance: createControlPlaneLogger(config, loggerDestination),
logController: new LogController({ disableRequestLogging: true }),
genReqId: (rawRequest) => normalizeRequestId(rawRequest.headers['x-request-id']),
trustProxy: true,
bodyLimit: Math.min(75_000_000, Math.max(2_000_000, Math.ceil(config.ARTIFACT_MAX_BYTES * 1.4) + 1_000_000))
});
@@ -271,6 +341,52 @@ export async function buildServer({
prefix: '/'
});
app.log.info({
diagnostic_event: 'service.initialized',
diagnostic_stage: 'startup',
host: config.HOST,
port: config.PORT,
log_level: config.LOG_LEVEL,
database_schema: config.DATABASE_SCHEMA,
database_ssl: config.DATABASE_SSL,
artifact_storage_backend: config.ARTIFACT_STORAGE_BACKEND,
agentbus_enabled: config.agentBusEnabled,
parser_loop_enabled: startParserLoop,
data_retention_enabled: config.DATA_RETENTION_ENABLED,
raw_payload_logging: config.AGENTBUS_LOG_PAYLOADS
}, 'control plane initialized');
app.addHook('onRequest', async (request, reply) => {
requestStartedAt.set(request, process.hrtime.bigint());
reply.header('x-request-id', requestId(request));
const contentLength = Number(request.headers['content-length']);
request.log.info({
diagnostic_event: 'http.request.started',
diagnostic_stage: 'http',
request_id: requestId(request),
method: request.method,
path: diagnosticRequestPath(request.url),
...(Number.isSafeInteger(contentLength) && contentLength >= 0 ? { content_length: contentLength } : {})
}, 'HTTP request started');
});
app.addHook('onResponse', async (request, reply) => {
const startedAt = requestStartedAt.get(request);
requestStartedAt.delete(request);
const metadata = {
diagnostic_event: 'http.request.completed',
diagnostic_stage: 'http',
request_id: requestId(request),
method: request.method,
path: diagnosticRequestPath(request.routeOptions.url || request.url),
status_code: reply.statusCode,
...(startedAt ? { duration_ms: diagnosticDurationMs(startedAt) } : {})
};
if (reply.statusCode >= 500) request.log.error(metadata, 'HTTP request completed with server error');
else if (reply.statusCode >= 400) request.log.warn(metadata, 'HTTP request completed with client error');
else request.log.info(metadata, 'HTTP request completed');
});
app.get('/history', async (_request, reply) => {
reply.header('Cache-Control', 'no-store');
return reply.sendFile('index.html');
@@ -292,8 +408,16 @@ export async function buildServer({
if (redirectUrl) return reply.redirect(redirectUrl, 308);
});
const auth = new AuthService(config);
const tasks = new TaskService(config);
const auth = new AuthService(config, {
info: (metadata, message) => app.log.info(metadata, message),
warn: (metadata, message) => app.log.warn(metadata, message),
error: (metadata, message) => app.log.error(metadata, message)
});
const tasks = new TaskService(config, undefined, {
info: (metadata, message) => app.log.info(metadata, message),
warn: (metadata, message) => app.log.warn(metadata, message),
error: (metadata, message) => app.log.error(metadata, message)
});
const externalParser = parser || await loadExternalParser();
const parserOrchestrator = new ParserOrchestrator(externalParser);
const activeParseWorkers = new Map<string, string>();
@@ -302,7 +426,11 @@ export async function buildServer({
let aiProbeValue: unknown;
let aiProbeExpiresAt = 0;
let aiProbeInFlight: Promise<unknown> | null = null;
const channelService = new AgentBusChannelService(config);
const channelService = new AgentBusChannelService(config, {
info: (metadata, message) => app.log.info(agentBusDiagnosticMetadata(metadata), message),
warn: (metadata, message) => app.log.warn(agentBusDiagnosticMetadata(metadata), message),
error: (metadata, message) => app.log.error(agentBusDiagnosticMetadata(metadata), message)
});
let agentBus: AgentBusManager | null = null;
const getSession = async (request: FastifyRequest): Promise<ActiveSession> => {
@@ -333,6 +461,7 @@ export async function buildServer({
result: unknown,
decision?: ParseDecisionInput
): Promise<void> {
const startedAt = process.hrtime.bigint();
const context: TaskContext = {
organizationId: claim.task.organization_id,
userId: '',
@@ -350,29 +479,65 @@ export async function buildServer({
decision
);
app.log.info({
diagnostic_event: 'parser.outcome.persisted',
diagnostic_stage: 'parser_persistence',
request_id: context.requestId,
task_id: finalized.task_id,
attempt_no: claim.attemptNo,
status: finalized.status
status: finalized.status,
duration_ms: diagnosticDurationMs(startedAt)
}, 'parse task finalized');
return;
} catch (error) {
if (error instanceof TaskError && ['stale_parse_result', 'task_cancelled', 'task_not_found'].includes(error.code)) {
app.log.info({ task_id: claim.task.task_id, attempt_no: claim.attemptNo, error_code: error.code }, 'late parse outcome ignored');
app.log.info({
diagnostic_event: 'parser.outcome.ignored',
diagnostic_stage: 'parser_persistence',
request_id: context.requestId,
task_id: claim.task.task_id,
attempt_no: claim.attemptNo,
error_code: error.code,
duration_ms: diagnosticDurationMs(startedAt)
}, 'late parse outcome ignored');
return;
}
lastError = error;
app.log.warn({
diagnostic_event: 'parser.outcome.persist_retry',
diagnostic_stage: 'parser_persistence',
request_id: context.requestId,
task_id: claim.task.task_id,
attempt_no: claim.attemptNo,
persistence_attempt: attempt + 1,
retrying: attempt < 2,
...diagnosticError(error, workerErrorCode(error))
}, 'parse outcome persistence attempt failed');
if (attempt < 2) await waitMs(250 * (attempt + 1));
}
}
app.log.error({
diagnostic_event: 'parser.outcome.persist_failed',
diagnostic_stage: 'parser_persistence',
request_id: context.requestId,
task_id: claim.task.task_id,
attempt_no: claim.attemptNo,
error_code: workerErrorCode(lastError),
error_name: lastError instanceof Error ? lastError.name : 'unknown'
duration_ms: diagnosticDurationMs(startedAt),
...diagnosticError(lastError, workerErrorCode(lastError))
}, 'parse outcome persistence failed after retries');
}
async function runParseTask(claim: ParseTaskClaim): Promise<void> {
const startedAt = process.hrtime.bigint();
const parseRequestId = `parse:${claim.task.task_id}:attempt:${claim.attemptNo}`;
app.log.info({
diagnostic_event: 'parser.worker.started',
diagnostic_stage: 'parser_execution',
request_id: parseRequestId,
task_id: claim.task.task_id,
attempt_no: claim.attemptNo,
parser_mode: claim.task.parser.configured_mode,
business_route_id: claim.task.parser.route_id
}, 'parse worker started');
const controller = new AbortController();
let timeout: NodeJS.Timeout | undefined;
const parserPromise = Promise.resolve().then(() => parserOrchestrator.parse(claim, controller.signal));
@@ -391,6 +556,16 @@ export async function buildServer({
})
]);
if (outcome.timedOut) {
app.log.error({
diagnostic_event: 'parser.worker.timeout',
diagnostic_stage: 'parser_execution',
request_id: parseRequestId,
task_id: claim.task.task_id,
attempt_no: claim.attemptNo,
timeout_ms: config.PARSE_WORKER_TIMEOUT_MS,
duration_ms: diagnosticDurationMs(startedAt),
error_code: 'parse_worker_timeout'
}, 'parse worker timed out');
await persistParseOutcome(
claim,
workerFailureResult('parse_worker_timeout', 'worker_watchdog_timeout')
@@ -398,13 +573,24 @@ export async function buildServer({
return;
}
await persistParseOutcome(claim, outcome.result.result, outcome.result.decision);
app.log.info({
diagnostic_event: 'parser.worker.completed',
diagnostic_stage: 'parser_execution',
request_id: parseRequestId,
task_id: claim.task.task_id,
attempt_no: claim.attemptNo,
duration_ms: diagnosticDurationMs(startedAt)
}, 'parse worker completed');
} catch (error) {
const errorCode = workerErrorCode(error);
app.log.error({
diagnostic_event: 'parser.worker.failed',
diagnostic_stage: 'parser_execution',
request_id: parseRequestId,
task_id: claim.task.task_id,
attempt_no: claim.attemptNo,
error_code: errorCode,
error_name: error instanceof Error ? error.name : 'unknown'
duration_ms: diagnosticDurationMs(startedAt),
...diagnosticError(error, errorCode)
}, 'parse worker failed');
await persistParseOutcome(claim, workerFailureResult(errorCode, 'worker_exception'));
} finally {
@@ -415,28 +601,48 @@ export async function buildServer({
async function processParseQueue(): Promise<void> {
const recovered = await tasks.recoverExpiredParseTasks();
if (recovered > 0) {
app.log.warn({ recovered_tasks: recovered }, 'expired parse tasks were durably blocked');
app.log.warn({
diagnostic_event: 'parser.queue.expired_tasks_recovered',
diagnostic_stage: 'parser_queue',
recovered_tasks: recovered
}, 'expired parse tasks were durably blocked');
}
const recoveredExecutions = await tasks.recoverExpiredExecutionTasks();
if (recoveredExecutions > 0) {
app.log.warn({ recovered_tasks: recoveredExecutions }, 'expired ERP executions were automatically failed and released');
app.log.warn({
diagnostic_event: 'erp.execution.expired_tasks_recovered',
diagnostic_stage: 'erp_execution',
recovered_tasks: recoveredExecutions
}, 'expired ERP executions were automatically failed and released');
}
const staleReconciliations = await tasks.maintainStaleReconciliationTasks();
if (staleReconciliations > 0) {
app.log.warn({ stale_tasks: staleReconciliations }, 'stale reconciliation tasks were automatically failed and released');
app.log.warn({
diagnostic_event: 'erp.reconciliation.stale_tasks_recovered',
diagnostic_stage: 'erp_reconciliation',
stale_tasks: staleReconciliations
}, 'stale reconciliation tasks were automatically failed and released');
}
while (activeParseWorkers.size < 2) {
const workerId = `parse:${process.pid}:${randomUUID()}`;
const claim = await tasks.claimNextParseTask(workerId, [...activeParseWorkers.keys()]);
if (!claim) break;
activeParseWorkers.set(claim.task.task_id, workerId);
app.log.info({
diagnostic_event: 'parser.queue.claimed',
diagnostic_stage: 'parser_queue',
task_id: claim.task.task_id,
attempt_no: claim.attemptNo,
active_workers: activeParseWorkers.size
}, 'parse queue task claimed');
void runParseTask(claim)
.catch((error) => {
app.log.error({
diagnostic_event: 'parser.runner.crashed',
diagnostic_stage: 'parser_execution',
task_id: claim.task.task_id,
attempt_no: claim.attemptNo,
error_code: workerErrorCode(error),
error_name: error instanceof Error ? error.name : 'unknown'
...diagnosticError(error, workerErrorCode(error))
}, 'parse task runner crashed');
})
.finally(() => {
@@ -454,7 +660,11 @@ export async function buildServer({
const now = Date.now();
if (now - lastParseQueueErrorAt < 60_000) return;
lastParseQueueErrorAt = now;
app.log.error({ error_code: workerErrorCode(error), error_name: error instanceof Error ? error.name : 'unknown' }, 'parse queue tick failed');
app.log.error({
diagnostic_event: 'parser.queue.tick_failed',
diagnostic_stage: 'parser_queue',
...diagnosticError(error, workerErrorCode(error))
}, 'parse queue tick failed');
})
.finally(() => {
parseQueueInFlight = null;
@@ -476,9 +686,9 @@ export async function buildServer({
organizationId: organization.id,
scheduleParseQueue,
logger: {
info: (metadata, message) => app.log.info(metadata, message),
warn: (metadata, message) => app.log.warn(metadata, message),
error: (metadata, message) => app.log.error(metadata, message)
info: (metadata, message) => app.log.info(agentBusDiagnosticMetadata(metadata), message),
warn: (metadata, message) => app.log.warn(agentBusDiagnosticMetadata(metadata), message),
error: (metadata, message) => app.log.error(agentBusDiagnosticMetadata(metadata), message)
}
});
await agentBus.start();
@@ -496,7 +706,11 @@ export async function buildServer({
try {
return await externalParser.checkConnection();
} catch (error) {
app.log.warn({ error: error instanceof Error ? error.message : String(error) }, 'external parser status probe failed');
app.log.warn({
diagnostic_event: 'parser.status_probe.failed',
diagnostic_stage: 'parser_probe',
...diagnosticError(error, 'probe_failed')
}, 'external parser status probe failed');
return { configured: Boolean(config.DEERFLOW_OPEN_API_KEY), reachable: false, ok: false, error_code: 'probe_failed' };
}
})();
@@ -513,6 +727,14 @@ export async function buildServer({
app.get('/health/ready', async (_request, reply) => {
const readiness = await databaseReadiness(config);
if (!readiness.ready) {
app.log.warn({
diagnostic_event: 'health.readiness.failed',
diagnostic_stage: 'health',
database_ready: readiness.database,
schema_ready: readiness.schema,
required_migration: readiness.requiredMigration,
error_code: readiness.errorCode || (readiness.database ? 'database_schema_outdated' : 'database_unavailable')
}, 'control plane readiness check failed');
return reply.code(503).send({
ok: false,
database: readiness.database,
@@ -521,6 +743,14 @@ export async function buildServer({
...(readiness.errorCode ? { error_code: readiness.errorCode } : {})
});
}
app.log.info({
diagnostic_event: 'health.readiness.passed',
diagnostic_stage: 'health',
database_ready: true,
schema_ready: true,
required_migration: readiness.requiredMigration,
agentbus_connected: agentBus?.status().connected || false
}, 'control plane readiness check passed');
return {
ok: true,
database: true,
@@ -909,6 +1139,16 @@ export async function buildServer({
app.setErrorHandler((error, request, reply) => {
if (error instanceof AuthError || error instanceof TaskError) {
request.log.warn({
diagnostic_event: 'http.request.rejected',
diagnostic_stage: 'http',
request_id: requestId(request),
...(requestTaskId(request) ? { task_id: requestTaskId(request) } : {}),
status_code: error.statusCode,
error_code: error.code,
error_name: error.name,
...(error instanceof TaskError ? { error_detail_keys: diagnosticMetadataKeys(error.details) } : {})
}, 'HTTP request rejected by application guard');
return reply.code(error.statusCode).send({
ok: false,
error_code: error.code,
@@ -917,9 +1157,25 @@ export async function buildServer({
});
}
if (error instanceof z.ZodError) {
request.log.warn({
diagnostic_event: 'http.request.invalid',
diagnostic_stage: 'http',
request_id: requestId(request),
...(requestTaskId(request) ? { task_id: requestTaskId(request) } : {}),
status_code: 400,
error_code: 'invalid_request',
validation_paths: error.issues.map((issue) => issue.path.join('.')).slice(0, 100)
}, 'HTTP request validation failed');
return reply.code(400).send({ ok: false, error_code: 'invalid_request', message: '请求参数不符合要求。', details: error.issues.map((issue) => issue.path.join('.')) });
}
request.log.error({ error: error instanceof Error ? error.message : String(error), request_id: requestId(request) }, 'unhandled request error');
request.log.error({
diagnostic_event: 'http.request.failed',
diagnostic_stage: 'http',
request_id: requestId(request),
...(requestTaskId(request) ? { task_id: requestTaskId(request) } : {}),
status_code: 500,
...diagnosticError(error, 'server_error')
}, 'unhandled request error');
return reply.code(500).send({ ok: false, error_code: 'server_error', message: '服务暂时不可用。' });
});
@@ -928,18 +1184,75 @@ export async function buildServer({
app.addHook('onClose', async () => clearInterval(interval));
}
app.addHook('onClose', async () => closePool());
app.addHook('onClose', async () => {
app.log.info({
diagnostic_event: 'service.closing',
diagnostic_stage: 'shutdown'
}, 'control plane closing');
await closePool();
});
return { app, auth, tasks, agentBus, channelService };
}
function installProcessDiagnostics(app: Awaited<ReturnType<typeof buildServer>>['app']): void {
let shuttingDown = false;
const shutdown = async (reason: string, exitCode: number, error?: unknown): Promise<void> => {
if (shuttingDown) return;
shuttingDown = true;
process.exitCode = exitCode;
const metadata = {
diagnostic_event: 'service.shutdown.started',
diagnostic_stage: 'shutdown',
shutdown_reason: reason,
exit_code: exitCode,
...(error === undefined ? {} : diagnosticError(error, reason))
};
if (exitCode === 0) app.log.info(metadata, 'control plane shutdown started');
else app.log.error(metadata, 'control plane shutdown started after fatal error');
try {
await app.close();
app.log.info({
diagnostic_event: 'service.shutdown.completed',
diagnostic_stage: 'shutdown',
shutdown_reason: reason,
exit_code: exitCode
}, 'control plane shutdown completed');
} catch (closeError) {
app.log.error({
diagnostic_event: 'service.shutdown.failed',
diagnostic_stage: 'shutdown',
shutdown_reason: reason,
exit_code: 1,
...diagnosticError(closeError, 'shutdown_failed')
}, 'control plane shutdown failed');
process.exitCode = 1;
}
};
process.once('SIGTERM', () => void shutdown('sigterm', 0));
process.once('SIGINT', () => void shutdown('sigint', 0));
process.once('uncaughtException', (error) => void shutdown('uncaught_exception', 1, error));
process.once('unhandledRejection', (error) => void shutdown('unhandled_rejection', 1, error));
}
async function main(): Promise<void> {
const config = loadConfig();
await assertDatabaseSchema(config);
const { app } = await buildServer({ config });
installProcessDiagnostics(app);
try {
await app.listen({ host: config.HOST, port: config.PORT });
app.log.info({ host: config.HOST, port: config.PORT }, 'LianSyn-platform control plane listening');
app.log.info({
diagnostic_event: 'service.listening',
diagnostic_stage: 'startup',
host: config.HOST,
port: config.PORT
}, 'LianSyn-platform control plane listening');
} catch (error) {
app.log.error({
diagnostic_event: 'service.listen.failed',
diagnostic_stage: 'startup',
...diagnosticError(error, 'listen_failed')
}, 'control plane failed to listen');
await app.close().catch(() => undefined);
throw error;
}
@@ -947,7 +1260,7 @@ async function main(): Promise<void> {
if (process.argv[1] && resolve(process.argv[1]) === resolve(fileURLToPath(import.meta.url))) {
main().catch(async (error) => {
console.error(error);
writeEmergencyDiagnostic('service.startup_failed', error, { diagnostic_stage: 'startup' });
await closePool();
process.exitCode = 1;
});