package com.xcong.excoin.modules.station.service.impl; import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSONObject; 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 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 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 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 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 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 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()); } } }