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