Administrator
2026-08-13 a37c520150f077649f7d8b314a517f09956ffae9
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
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());
        }
    }
}