Administrator
2026-08-14 57e5a8385de2d5ff2e43c72535b472e42d65ba91
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
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
package com.xcong.excoin.modules.station.service.impl;
 
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.xcong.excoin.modules.station.dao.StrategyEventLogDao;
import com.xcong.excoin.modules.station.dao.StrategyStatusDao;
import com.xcong.excoin.modules.station.entity.StrategyEventLog;
import com.xcong.excoin.modules.station.entity.StrategyStatus;
import com.xcong.excoin.modules.station.model.GateStatsEvent;
import com.xcong.excoin.modules.station.producer.CmdProducer;
import com.xcong.excoin.modules.station.service.StationService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
 
import javax.annotation.Resource;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
import java.nio.charset.StandardCharsets;
import java.util.Date;
import java.util.List;
import java.util.Map;
 
@Slf4j
@Service
public class StationServiceImpl extends ServiceImpl<StrategyEventLogDao, StrategyEventLog> implements StationService {
 
    @Resource
    private CmdProducer cmdProducer;
 
    @Resource
    private StrategyStatusDao strategyStatusDao;
 
    @Resource
    private StrategyEventLogDao strategyEventLogDao;
 
    // ==================== 启停 ====================
 
    @Override
    public void startStrategy(String apiKey) {
        cmdProducer.sendCommand(md5(apiKey), "START");
    }
 
    @Override
    public void stopStrategy(String apiKey) {
        cmdProducer.sendCommand(md5(apiKey), "STOP");
    }
 
    @Override
    public void startByMd5(String apiKeyMd5) {
        cmdProducer.sendCommand(apiKeyMd5, "START");
    }
 
    @Override
    public void stopByMd5(String apiKeyMd5) {
        cmdProducer.sendCommand(apiKeyMd5, "STOP");
    }
 
    @Override
    public void updateConfig(String apiKeyMd5, Map<String, Object> params) {
        // 过滤空字符串,避免 Fastjson 反序列化到 BigDecimal/int 时抛异常
        params.entrySet().removeIf(e -> e.getValue() == null || "".equals(e.getValue()));
        // 1. 持久化配置到 strategy_status,保证保存后刷新页面可回显
        persistConfig(apiKeyMd5, params);
        // 2. 发送 MQ 命令,通知 JAR(需手动点击「启动」才生效)
        cmdProducer.sendCommand(apiKeyMd5, "UPDATE_CONFIG", JSON.toJSONString(params));
    }
 
    // ==================== 别名 ====================
 
    @Override
    public void saveAlias(String apiKeyMd5, String alias) {
        String cleaned = alias == null ? "" : alias.trim();
        String value = cleaned.isEmpty() ? null : cleaned;
 
        StrategyStatus existing = strategyStatusDao.selectByApiKeyMd5(apiKeyMd5);
        if (existing != null) {
            // 用 UpdateWrapper 显式 set,支持把别名更新为 null(清除)
            LambdaUpdateWrapper<StrategyStatus> wrapper = new LambdaUpdateWrapper<>();
            wrapper.eq(StrategyStatus::getApiKeyMd5, apiKeyMd5)
                   .set(StrategyStatus::getAliasName, value);
            strategyStatusDao.update(null, wrapper);
        } else if (value != null) {
            // 实例尚无状态记录(未启动过),仅创建一条只含别名的记录
            StrategyStatus status = new StrategyStatus();
            status.setApiKeyMd5(apiKeyMd5);
            status.setAliasName(value);
            status.setCreateTime(new Date());
            strategyStatusDao.insert(status);
        }
    }
 
    @Override
    public String getAlias(String apiKeyMd5) {
        StrategyStatus status = strategyStatusDao.selectByApiKeyMd5(apiKeyMd5);
        return status == null ? null : status.getAliasName();
    }
 
    // ==================== 事件落库 ====================
 
    @Override
    public void saveEvent(GateStatsEvent event) {
        // 幂等去重
        if (strategyEventLogDao.countByEventId(event.getEventId()) > 0) {
            return;
        }
 
        StrategyEventLog log = new StrategyEventLog();
        log.setEventId(event.getEventId());
        log.setEventType(event.getType());
        log.setApiKeyMd5(event.getApiKeyMd5());
        log.setEventTime(event.getTimestamp());
        log.setPayloadJson(event.getPayload());
        log.setCreateTime(new Date());
 
        // 从 payload 提取 contract(仅 STRATEGY_START 时有)
        String contract = extractContract(event);
        log.setContract(contract);
 
        strategyEventLogDao.insert(log);
 
        // STRATEGY_START / STRATEGY_STOP 事件同步更新实时状态表
        if ("STRATEGY_START".equals(event.getType())) {
            upsertStrategyStatus(event, contract);
        } else if ("STRATEGY_STOP".equals(event.getType())) {
            updateStrategyStopped(event);
        }
    }
 
    @Override
    public List<StrategyEventLog> getEvents(String apiKeyMd5) {
        return strategyEventLogDao.selectByApiKeyMd5(apiKeyMd5);
    }
 
    @Override
    public StrategyStatus getStrategyStatus(String apiKeyMd5) {
        return strategyStatusDao.selectByApiKeyMd5(apiKeyMd5);
    }
 
    // ==================== 辅助方法 ====================
 
    private void upsertStrategyStatus(GateStatsEvent event, String contract) {
        JSONObject payload = JSON.parseObject(event.getPayload());
        StrategyStatus existing = strategyStatusDao.selectByApiKeyMd5(event.getApiKeyMd5());
 
        StrategyStatus status = (existing != null) ? existing : new StrategyStatus();
        status.setApiKeyMd5(event.getApiKeyMd5());
        status.setContract(contract);
        status.setState("ACTIVE");
        // 复用统一的配置字段映射,保证 START 事件落库字段与 updateConfig 一致
        applyConfigFields(status, payload);
 
        if (existing != null) {
            strategyStatusDao.updateById(status);
        } else {
            status.setCreateTime(new Date());
            strategyStatusDao.insert(status);
        }
    }
 
    private void updateStrategyStopped(GateStatsEvent event) {
        JSONObject payload = JSON.parseObject(event.getPayload());
        StrategyStatus status = strategyStatusDao.selectByApiKeyMd5(event.getApiKeyMd5());
        if (status != null) {
            status.setState("STOPPED");
            status.setCumulativePnl(payload.getString("pnl"));
            strategyStatusDao.updateById(status);
        }
    }
 
    /**
     * 持久化配置到 strategy_status(UPSERT)。
     * 保存配置但策略未启动时,DB 里只有配置字段,state/contract 保持原值或为空,
     * 前端刷新即可回显已保存参数。
     */
    private void persistConfig(String apiKeyMd5, Map<String, Object> params) {
        StrategyStatus existing = strategyStatusDao.selectByApiKeyMd5(apiKeyMd5);
        StrategyStatus status = (existing != null) ? existing : new StrategyStatus();
        status.setApiKeyMd5(apiKeyMd5);
        applyConfigFields(status, params);
 
        if (existing != null) {
            strategyStatusDao.updateById(status);
        } else {
            status.setCreateTime(new Date());
            strategyStatusDao.insert(status);
        }
    }
 
    /**
     * 将配置 Map 的字段映射到 StrategyStatus 实体。
     * 前端提交的 key(gridRate/quantity/expectedProfit/rounds/priceDriveEnabled...)
     * 与 STRATEGY_START payload 的 key 一致,rounds 映射到实体字段 totalRounds。
     */
    private void applyConfigFields(StrategyStatus status, Map<String, Object> params) {
        if (params.containsKey("leverage")) status.setLeverage(str(params.get("leverage")));
        if (params.containsKey("contract")) status.setContract(str(params.get("contract")));
        if (params.containsKey("gridRate")) status.setGridRate(str(params.get("gridRate")));
        if (params.containsKey("expectedProfit")) status.setExpectedProfit(str(params.get("expectedProfit")));
        if (params.containsKey("maxLoss")) status.setMaxLoss(str(params.get("maxLoss")));
        if (params.containsKey("baseQuantity")) status.setBaseQuantity(str(params.get("baseQuantity")));
        if (params.containsKey("quantity")) status.setQuantity(str(params.get("quantity")));
        if (params.containsKey("maxPositionSize")) status.setMaxPositionSize(intVal(params.get("maxPositionSize")));
        if (params.containsKey("stopLossCount")) status.setStopLossCount(intVal(params.get("stopLossCount")));
        if (params.containsKey("takeProfitGridSpan")) status.setTakeProfitGridSpan(intVal(params.get("takeProfitGridSpan")));
        if (params.containsKey("stopLossCountMode")) status.setStopLossCountMode(str(params.get("stopLossCountMode")));
        if (params.containsKey("addPositionInterval")) status.setAddPositionInterval(intVal(params.get("addPositionInterval")));
        if (params.containsKey("addPositionQuantity")) status.setAddPositionQuantity(intVal(params.get("addPositionQuantity")));
        if (params.containsKey("maxPositionPerSide")) status.setMaxPositionPerSide(intVal(params.get("maxPositionPerSide")));
        if (params.containsKey("addPositionStartThreshold")) status.setAddPositionStartThreshold(intVal(params.get("addPositionStartThreshold")));
        if (params.containsKey("placeExcessTakeProfit")) status.setPlaceExcessTakeProfit(boolVal(params.get("placeExcessTakeProfit")));
        if (params.containsKey("priceDriveEnabled")) status.setPriceDriveEnabled(boolVal(params.get("priceDriveEnabled")));
        if (params.containsKey("rounds")) status.setTotalRounds(intVal(params.get("rounds")));
        if (params.containsKey("principal")) status.setPrincipal(str(params.get("principal")));
    }
 
    private static String str(Object o) {
        return o == null ? null : String.valueOf(o);
    }
 
    private static Integer intVal(Object o) {
        if (o == null) return null;
        if (o instanceof Number) return ((Number) o).intValue();
        try {
            return Integer.parseInt(String.valueOf(o).trim());
        } catch (NumberFormatException e) {
            return null;
        }
    }
 
    private static Boolean boolVal(Object o) {
        if (o == null) return null;
        if (o instanceof Boolean) return (Boolean) o;
        String s = String.valueOf(o).trim();
        return "true".equalsIgnoreCase(s) || "1".equals(s);
    }
 
    private String extractContract(GateStatsEvent event) {
        try {
            JSONObject payload = JSON.parseObject(event.getPayload());
            return payload.getString("contract");
        } catch (Exception e) {
            return null;
        }
    }
 
    private static String md5(String input) {
        try {
            MessageDigest md = MessageDigest.getInstance("MD5");
            byte[] digest = md.digest(input.getBytes(StandardCharsets.UTF_8));
            StringBuilder sb = new StringBuilder();
            for (byte b : digest) sb.append(String.format("%02x", b));
            return sb.toString();
        } catch (NoSuchAlgorithmException e) {
            return Integer.toHexString(input.hashCode());
        }
    }
}