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