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()); } } }