From 747a86bb94006aaca721cfc0c0ce7061643a9ea6 Mon Sep 17 00:00:00 2001
From: chenyc <501753378@qq.com>
Date: 星期四, 20 八月 2026 12:19:35 +0800
Subject: [PATCH] gx
---
app.js | 107 +++++++++++++++++++++++++++++++++--------------------
1 files changed, 67 insertions(+), 40 deletions(-)
diff --git a/app.js b/app.js
index 7c4f879..124f0e9 100644
--- a/app.js
+++ b/app.js
@@ -15,6 +15,7 @@
error: 40,
};
+// 统一生成本地日志时间,兼顾控制台和文件日志格式。
function formatLocalTimestamp(date = new Date()) {
const year = date.getFullYear();
const month = String(date.getMonth() + 1).padStart(2, '0');
@@ -101,7 +102,7 @@
try {
fs.appendFileSync(logFilePath, `${line}\n`, 'utf8');
} catch (error) {
- console.error(`[LOGGER] Failed to append log file: ${error.message}`);
+ console.error(`[LOGGER] 写入日志文件失败: ${error.message}`);
}
}
}
@@ -150,23 +151,23 @@
}
function printHelp() {
- console.log(`JHM TCP Socket Gateway
+ console.log(`JHM TCP Socket 网关服务
-Usage:
+用法:
node app.js
node app.js --config ./config.json
jhm-service.exe --config .\\runtime\\config.json
./jhm-service --config ./runtime/config.json
-Options:
- --config <path> set config file path
- --help, -h show help
+参数:
+ --config <path> 指定配置文件路径
+ --help, -h 显示帮助
`);
}
function loadConfig(configPath) {
if (!fs.existsSync(configPath)) {
- throw new Error(`config file not found: ${configPath}`);
+ throw new Error(`未找到配置文件: ${configPath}`);
}
const content = fs.readFileSync(configPath, 'utf8');
@@ -208,66 +209,67 @@
return {
publishTime: bloodPressure ? bloodPressure.publishTime !== false : true,
+ flushImmediately: bloodPressure ? bloodPressure.flushImmediately !== false : true,
};
}
function validateConfig(config) {
if (!config.tcp || !Array.isArray(config.devices)) {
- throw new Error('config.json must include tcp and devices');
+ throw new Error('config.json 必须包含 tcp 和 devices');
}
if (!config.tcp.host || !config.tcp.port) {
- throw new Error('config.json must include tcp.host and tcp.port');
+ throw new Error('config.json 必须包含 tcp.host 和 tcp.port');
}
const channels = getSendChannels(config);
if (channels.length === 0) {
- throw new Error('config.json send.channels must enable at least one channel');
+ throw new Error('config.json 的 send.channels 至少要启用一个通道');
}
if (channels.includes('mqtt')) {
if (!config.mqtt) {
- throw new Error('config.json enabled mqtt but missing mqtt config');
+ throw new Error('config.json 已启用 mqtt,但缺少 mqtt 配置');
}
const hasBrokerUrl = Boolean(config.mqtt.brokerUrl);
const hasHostMode = Boolean(config.mqtt.protocol && config.mqtt.host && config.mqtt.port);
if (!hasBrokerUrl && !hasHostMode) {
- throw new Error('config.json mqtt requires brokerUrl or protocol/host/port');
+ throw new Error('config.json 的 mqtt 必须配置 brokerUrl 或 protocol/host/port');
}
if (!config.mqtt.topicTemplate && !config.mqtt.defaultTopicPrefix) {
- throw new Error('config.json mqtt requires topicTemplate or defaultTopicPrefix');
+ throw new Error('config.json 的 mqtt 必须配置 topicTemplate 或 defaultTopicPrefix');
}
}
if (channels.includes('aliyun')) {
if (!config.aliyun) {
- throw new Error('config.json enabled aliyun but missing aliyun config');
+ throw new Error('config.json 已启用 aliyun,但缺少 aliyun 配置');
}
if (!config.aliyun.tupleApiBaseUrl && !config.aliyun.tupleApiUrl) {
- throw new Error('config.json aliyun requires tupleApiBaseUrl or tupleApiUrl');
+ throw new Error('config.json 的 aliyun 必须配置 tupleApiBaseUrl 或 tupleApiUrl');
}
}
if (!config.protocol || !config.protocol.alModelPath) {
- throw new Error('config.json protocol.alModelPath is required');
+ throw new Error('config.json 必须配置 protocol.alModelPath');
}
const sendOptions = getSendOptions(config);
if (sendOptions.flushIntervalMs <= 0) {
- throw new Error('config.json send.flushIntervalMs must be > 0');
+ throw new Error('config.json 的 send.flushIntervalMs 必须大于 0');
}
if (config.logging && config.logging.level) {
const supportedLevels = ['debug', 'info', 'warn', 'error'];
if (!supportedLevels.includes(String(config.logging.level).toLowerCase())) {
- throw new Error(`config.json logging.level only supports ${supportedLevels.join(', ')}`);
+ throw new Error(`config.json 的 logging.level 仅支持 ${supportedLevels.join('、')}`);
}
}
}
@@ -278,6 +280,29 @@
}
return path.join(path.dirname(configFilePath), targetPath);
+}
+
+function createOnMetricHandler({ logger, aggregator, sendOptions, bloodPressureOptions }) {
+ return function onMetric(device, metric, result) {
+ logger.info(`[APP] 指标已缓存 deviceId=${device.deviceId} mode=${sendOptions.mode} metric=${JSON.stringify(metric)}`);
+ return aggregator.ingest(device, metric)
+ .then(() => {
+ if (
+ sendOptions.mode === 'batch'
+ && bloodPressureOptions.flushImmediately
+ && result
+ && result.protocol === 'blood-pressure'
+ ) {
+ logger.info(`[APP] 血压报文触发整包立即发送 deviceId=${device.deviceId}`);
+ return aggregator.flush({ reason: 'blood-pressure', deviceId: device.deviceId });
+ }
+
+ return null;
+ })
+ .catch((error) => {
+ logger.error(`[APP] 指标处理失败 deviceId=${device.deviceId}: ${error.message}`);
+ });
+ };
}
async function main() {
@@ -292,7 +317,7 @@
enabled: false,
console: true,
});
- bootstrapLogger.info(`[APP] Startup arguments parsed config=${options.configPath}`);
+ bootstrapLogger.info(`[APP] 启动参数解析完成 config=${options.configPath}`);
const config = loadConfig(options.configPath);
const logDir = resolveConfigPath(options.configPath, (config.logging && config.logging.dir) || './logs');
@@ -306,21 +331,22 @@
const alModelPath = resolveConfigPath(options.configPath, config.protocol.alModelPath);
if (config.logging && config.logging.enabled === false) {
- logger.info('[APP] File logging is disabled; console output only');
+ logger.info('[APP] 文件日志已禁用,仅输出到控制台');
} else {
- logger.info(`[APP] Local logging enabled dir=${logDir} file=${logger.getLogFilePath()} level=${normalizeLogLevel(config.logging && config.logging.level)}`);
+ logger.info(`[APP] 本地日志已启用 dir=${logDir} file=${logger.getLogFilePath()} level=${normalizeLogLevel(config.logging && config.logging.level)}`);
}
- logger.info(`[APP] Config loaded tcp=${config.tcp.host}:${config.tcp.port} devices=${config.devices.length} channels=${sendChannels.join(',')} sendMode=${sendOptions.mode} flushIntervalMs=${sendOptions.flushIntervalMs} alModel=${alModelPath}`);
+ logger.info(`[APP] 配置加载完成 tcp=${config.tcp.host}:${config.tcp.port} devices=${config.devices.length} channels=${sendChannels.join(',')} sendMode=${sendOptions.mode} flushIntervalMs=${sendOptions.flushIntervalMs} alModel=${alModelPath}`);
const mqttService = sendChannels.includes('mqtt') ? new MqttService(config.mqtt, logger) : null;
const aliyunService = sendChannels.includes('aliyun') ? new AliyunService(config.aliyun, logger) : null;
+ // 聚合器负责把多次单点指标合并成一次完整物模型上报。
const aggregator = new MetricAggregator({
logger,
...sendOptions,
onFlush: async (device, payload, meta) => {
- logger.info(`[APP] Batch dispatch deviceId=${device.deviceId} channels=${sendChannels.join(',')} reason=${meta.reason} payload=${JSON.stringify(payload)}`);
+ logger.info(`[APP] 开始批量分发 deviceId=${device.deviceId} channels=${sendChannels.join(',')} reason=${meta.reason} payload=${JSON.stringify(payload)}`);
let ok = true;
@@ -329,7 +355,7 @@
mqttService.publish(device, payload);
} catch (error) {
ok = false;
- logger.error(`[APP] MQTT batch dispatch failed deviceId=${device.deviceId}: ${error.message}`);
+ logger.error(`[APP] MQTT 批量分发失败 deviceId=${device.deviceId}: ${error.message}`);
}
}
@@ -339,14 +365,14 @@
if (result && result.skipped) {
ok = false;
- logger.warn(`[APP] Aliyun batch dispatch skipped deviceId=${device.deviceId} reason=${result.reason}`);
+ logger.warn(`[APP] 阿里云批量分发已跳过 deviceId=${device.deviceId} reason=${result.reason}`);
} else if (result && result.ok === false) {
ok = false;
- logger.error(`[APP] Aliyun batch dispatch failed deviceId=${device.deviceId}: ${result.reason}`);
+ logger.error(`[APP] 阿里云批量分发失败 deviceId=${device.deviceId}: ${result.reason}`);
}
} catch (error) {
ok = false;
- logger.error(`[APP] Aliyun batch dispatch exception deviceId=${device.deviceId}: ${error.message}`);
+ logger.error(`[APP] 阿里云批量分发异常 deviceId=${device.deviceId}: ${error.message}`);
}
}
@@ -360,27 +386,27 @@
alModelPath,
publishBloodPressureTime: bloodPressureOptions.publishTime,
logger,
- onMetric(device, metric) {
- logger.info(`[APP] Metric cached deviceId=${device.deviceId} mode=${sendOptions.mode} metric=${JSON.stringify(metric)}`);
- aggregator.ingest(device, metric).catch((error) => {
- logger.error(`[APP] Metric cache failed deviceId=${device.deviceId}: ${error.message}`);
- });
- },
+ onMetric: createOnMetricHandler({
+ logger,
+ aggregator,
+ sendOptions,
+ bloodPressureOptions,
+ }),
});
if (mqttService) {
- logger.info('[APP] Starting MQTT channel');
+ logger.info('[APP] 正在启动 MQTT 通道');
mqttService.start();
}
if (aliyunService) {
- logger.info('[APP] Starting Aliyun channel');
+ logger.info('[APP] 正在启动阿里云通道');
aliyunService.start();
}
aggregator.start();
await tcpService.start();
- logger.info('[APP] All services started');
+ logger.info('[APP] 所有服务已启动');
let shuttingDown = false;
@@ -390,7 +416,7 @@
}
shuttingDown = true;
- logger.warn(`[APP] Received ${signal}, shutting down`);
+ logger.warn(`[APP] 收到 ${signal},开始关闭服务`);
await tcpService.stop();
await aggregator.stop();
@@ -403,7 +429,7 @@
await aliyunService.stop();
}
- logger.info('[APP] Service shutdown completed');
+ logger.info('[APP] 服务已完成关闭');
await logger.close();
process.exit(0);
}
@@ -417,11 +443,11 @@
});
process.on('uncaughtException', (error) => {
- logger.error(`[APP] Uncaught exception: ${error.stack || error.message}`);
+ logger.error(`[APP] 未捕获异常: ${error.stack || error.message}`);
});
process.on('unhandledRejection', (reason) => {
- logger.error(`[APP] Unhandled promise rejection: ${reason}`);
+ logger.error(`[APP] 未处理的 Promise 拒绝: ${reason}`);
});
}
@@ -440,6 +466,7 @@
getSendOptions,
loadConfig,
main,
+ createOnMetricHandler,
normalizeLogLevel,
parseArgs,
printHelp,
--
Gitblit v1.8.0