From 392a6ba0f6c98b4a8b5398b4b75ad909f9bcf67c Mon Sep 17 00:00:00 2001
From: Administrator <15274802129@163.com>
Date: Mon, 24 Aug 2026 14:16:56 +0800
Subject: [PATCH] fix(gate): 修复心跳调度器中服务器端口获取问题
---
src/main/java/com/xcong/excoin/modules/gateApi/GateCommandConsumer.java | 111 +++++++++++++++++++++++++++++++++++++++++++++++++++++++
1 files changed, 111 insertions(+), 0 deletions(-)
diff --git a/src/main/java/com/xcong/excoin/modules/gateApi/GateCommandConsumer.java b/src/main/java/com/xcong/excoin/modules/gateApi/GateCommandConsumer.java
new file mode 100644
index 0000000..89907d8
--- /dev/null
+++ b/src/main/java/com/xcong/excoin/modules/gateApi/GateCommandConsumer.java
@@ -0,0 +1,111 @@
+package com.xcong.excoin.modules.gateApi;
+
+import com.alibaba.fastjson.JSON;
+import com.xcong.excoin.modules.station.model.CmdAckMsg;
+import com.xcong.excoin.modules.station.model.GateCommand;
+import com.xcong.excoin.modules.station.model.GateStatsEvent;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.amqp.rabbit.annotation.RabbitListener;
+import org.springframework.stereotype.Component;
+import org.springframework.util.StringUtils;
+
+import javax.annotation.Resource;
+
+/**
+ * JAR 侧 — 消费 Station 下发的启停指令
+ */
+@Slf4j
+@Component
+public class GateCommandConsumer {
+
+ @Resource
+ private StatsEventProducer statsProducer;
+
+ @Resource
+ private GateWebSocketClientManager manager;
+
+ /**
+ * 监听专属命令队列 QUEUE_GATE_CMD_{apiKeyMd5}
+ */
+ @RabbitListener(queues = "#{commandQueueInitializer.queueName}")
+ public void onCommand(String msg) {
+ GateCommand cmd;
+ try {
+ cmd = JSON.parseObject(msg, GateCommand.class);
+ } catch (Exception e) {
+ log.error("[Gate] 指令解析失败: {}", msg, e);
+ return;
+ }
+
+ String apiKeyMd5 = md5(manager.getConfig().getApiKey());
+ boolean success = true;
+ String newState = "";
+ String message = "";
+
+ GateGridTradeService strategy = manager.getGridTradeService();
+ try {
+ switch (cmd.getCommandType()) {
+ case "START":
+ manager.reloadAndStart();
+ newState = "ACTIVE";
+ message = "策略已启动(已加载最新配置)";
+ log.info("[Gate] 远程启动, 已重载配置");
+ break;
+
+ case "STOP":
+ strategy.stopGrid();
+ newState = strategy.getState().name();
+ message = "策略已停止";
+ log.info("[Gate] 远程停止, state={}", newState);
+ break;
+
+ case "UPDATE_CONFIG":
+ if (!StringUtils.hasText(cmd.getPayload())) {
+ success = false;
+ message = "UPDATE_CONFIG 缺少 payload";
+ log.warn("[Gate] UPDATE_CONFIG 指令缺少 payload");
+ break;
+ }
+ GateConfigDTO dto = JSON.parseObject(cmd.getPayload(), GateConfigDTO.class);
+ manager.updateConfig(dto);
+ newState = "ACTIVE";
+ message = "参数已持久化,待下次启动生效";
+ log.info("[Gate] 远程参数已持久化");
+ break;
+
+ default:
+ success = false;
+ message = "未知指令类型: " + cmd.getCommandType();
+ log.warn("[Gate] 未知指令: {}", cmd.getCommandType());
+ }
+ } catch (Exception e) {
+ success = false;
+ newState = "ERROR";
+ message = e.getMessage();
+ log.error("[Gate] 指令执行失败", e);
+ }
+
+ // 发送 ACK 到 Station
+ CmdAckMsg ack = CmdAckMsg.builder()
+ .commandId(cmd.getCommandId())
+ .success(success)
+ .newState(newState)
+ .message(message)
+ .build();
+ GateStatsEvent event = statsProducer.newCmdAck(apiKeyMd5, ack);
+ statsProducer.sendHeartbeat(event);
+ log.debug("[Gate] ACK 已发送, cmdId={}, success={}", cmd.getCommandId(), success);
+ }
+
+ private static String md5(String input) {
+ try {
+ java.security.MessageDigest md = java.security.MessageDigest.getInstance("MD5");
+ byte[] digest = md.digest(input.getBytes(java.nio.charset.StandardCharsets.UTF_8));
+ StringBuilder sb = new StringBuilder();
+ for (byte b : digest) sb.append(String.format("%02x", b));
+ return sb.toString();
+ } catch (Exception e) {
+ return Integer.toHexString(input.hashCode());
+ }
+ }
+}
--
Gitblit v1.9.1