1、反向下单 暂时提交
This commit is contained in:
@ -1,15 +0,0 @@
|
||||
package binanceservice
|
||||
|
||||
import "github.com/go-admin-team/go-admin-core/sdk/service"
|
||||
|
||||
type BinanceReverseFuturesService struct {
|
||||
service.Service
|
||||
}
|
||||
|
||||
//
|
||||
func (e *BinanceReverseFuturesService) DoReverse() error {
|
||||
apiInfo, ok := ShouldReverse(apiKey)
|
||||
//TODO: 实现反向开仓逻辑
|
||||
|
||||
return nil
|
||||
}
|
||||
@ -1,12 +0,0 @@
|
||||
package binanceservice
|
||||
|
||||
import "github.com/go-admin-team/go-admin-core/sdk/service"
|
||||
|
||||
type BinanceReverseSpotService struct {
|
||||
service.Service
|
||||
}
|
||||
|
||||
// JudgeApiUser 判断是否需要
|
||||
func (e *BinanceReverseSpotService) JudgeApiUser() bool {
|
||||
return false
|
||||
}
|
||||
@ -1,16 +0,0 @@
|
||||
package binanceservice
|
||||
|
||||
import DbModels "go-admin/app/admin/models"
|
||||
|
||||
// ShouldReverse 判断是否需要反单
|
||||
// return apiInfo, bool
|
||||
func ShouldReverse(apiKey string) (DbModels.LineApiUser, bool) {
|
||||
// TODO: 实现判断是否需要反单的逻辑
|
||||
apiInfo := GetApiInfoByKey(apiKey)
|
||||
|
||||
if apiInfo.ReverseStatus == 1 && apiInfo.ReverseApiId > 0 {
|
||||
return apiInfo, true
|
||||
}
|
||||
|
||||
return DbModels.LineApiUser{}, false
|
||||
}
|
||||
114
services/binanceservice/futures_reset_v2.go
Normal file
114
services/binanceservice/futures_reset_v2.go
Normal file
@ -0,0 +1,114 @@
|
||||
package binanceservice
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
DbModels "go-admin/app/admin/models"
|
||||
"go-admin/pkg/retryhelper"
|
||||
|
||||
"github.com/bytedance/sonic"
|
||||
log "github.com/go-admin-team/go-admin-core/logger"
|
||||
"github.com/go-admin-team/go-admin-core/sdk/service"
|
||||
)
|
||||
|
||||
type FuturesResetV2 struct {
|
||||
service.Service
|
||||
}
|
||||
|
||||
// 带重试机制的合约下单
|
||||
func (e *FuturesResetV2) OrderPlaceLoop(apiUserInfo *DbModels.LineApiUser, params FutOrderPlace) error {
|
||||
opts := retryhelper.DefaultRetryOptions()
|
||||
opts.RetryableErrFn = func(err error) bool {
|
||||
if strings.Contains(err.Error(), "LOT_SIZE") {
|
||||
return false
|
||||
}
|
||||
//重试
|
||||
return true
|
||||
}
|
||||
|
||||
err := retryhelper.Retry(func() error {
|
||||
return e.OrderPlace(apiUserInfo, params)
|
||||
}, opts)
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// 合约下单
|
||||
func (e *FuturesResetV2) OrderPlace(apiUserInfo *DbModels.LineApiUser, params FutOrderPlace) error {
|
||||
if err := params.CheckParams(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
orderType := strings.ToUpper(params.OrderType)
|
||||
side := strings.ToUpper(params.Side)
|
||||
|
||||
paramsMaps := map[string]string{
|
||||
"symbol": params.Symbol,
|
||||
"side": side,
|
||||
"quantity": params.Quantity.String(),
|
||||
"type": orderType,
|
||||
"newClientOrderId": params.NewClientOrderId,
|
||||
"positionSide": params.PositionSide,
|
||||
}
|
||||
|
||||
// 设置市价和限价等类型的参数
|
||||
switch orderType {
|
||||
case "LIMIT":
|
||||
paramsMaps["price"] = params.Price.String()
|
||||
paramsMaps["timeInForce"] = "GTC"
|
||||
case "TAKE_PROFIT_MARKET":
|
||||
paramsMaps["timeInForce"] = "GTC"
|
||||
paramsMaps["stopprice"] = params.Profit.String()
|
||||
paramsMaps["workingType"] = "MARK_PRICE"
|
||||
case "TAKE_PROFIT":
|
||||
paramsMaps["price"] = params.Price.String()
|
||||
paramsMaps["stopprice"] = params.Profit.String()
|
||||
paramsMaps["timeInForce"] = "GTC"
|
||||
case "STOP_MARKET":
|
||||
paramsMaps["stopprice"] = params.StopPrice.String()
|
||||
paramsMaps["workingType"] = "MARK_PRICE"
|
||||
paramsMaps["timeInForce"] = "GTC"
|
||||
case "STOP":
|
||||
paramsMaps["price"] = params.Price.String()
|
||||
paramsMaps["stopprice"] = params.StopPrice.String()
|
||||
paramsMaps["workingType"] = "MARK_PRICE"
|
||||
paramsMaps["timeInForce"] = "GTC"
|
||||
}
|
||||
|
||||
// 获取 API 信息和发送下单请求
|
||||
client := GetClient(apiUserInfo)
|
||||
_, statusCode, err := client.SendFuturesRequestAuth("/fapi/v1/order", "POST", paramsMaps)
|
||||
|
||||
if err != nil {
|
||||
return parseOrderError(err, paramsMaps, statusCode, apiUserInfo)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// 拆出的错误解析函数
|
||||
func parseOrderError(err error, paramsMaps map[string]string, statusCode int, apiUserInfo *DbModels.LineApiUser) error {
|
||||
var dataMap map[string]interface{}
|
||||
|
||||
if err2 := sonic.Unmarshal([]byte(err.Error()), &dataMap); err2 == nil {
|
||||
if code, ok := dataMap["code"]; ok {
|
||||
paramsVal, _ := sonic.MarshalString(¶msMaps)
|
||||
log.Error("下单失败 参数:", paramsVal)
|
||||
errContent := FutErrorMaps[code.(float64)]
|
||||
if errContent == "" {
|
||||
errContent = err.Error()
|
||||
}
|
||||
return fmt.Errorf("api_id:%d 交易对:%s 下单失败:%s", apiUserInfo.Id, paramsMaps["symbol"], errContent)
|
||||
}
|
||||
}
|
||||
|
||||
// 特殊处理保证金不足
|
||||
if strings.Contains(err.Error(), "Margin is insufficient.") {
|
||||
return fmt.Errorf("api_id:%d 交易对:%s 下单失败:%s", apiUserInfo.Id, paramsMaps["symbol"], FutErrorMaps[-2019])
|
||||
}
|
||||
|
||||
return fmt.Errorf("api_id:%d 交易对:%s statusCode:%v 下单失败:%s", apiUserInfo.Id, paramsMaps["symbol"], statusCode, err.Error())
|
||||
}
|
||||
@ -12,6 +12,7 @@ import (
|
||||
"go-admin/models/binancedto"
|
||||
"go-admin/models/futuresdto"
|
||||
"go-admin/pkg/httputils"
|
||||
"go-admin/pkg/retryhelper"
|
||||
"go-admin/pkg/utility"
|
||||
"strconv"
|
||||
"strings"
|
||||
@ -715,6 +716,15 @@ func (e FutRestApi) ClosePosition(symbol string, orderSn string, quantity decima
|
||||
return nil
|
||||
}
|
||||
|
||||
// 带重试机制的取消订单
|
||||
func (e FutRestApi) CancelFutOrderRetry(apiUserInfo DbModels.LineApiUser, symbol string, newClientOrderId string) error {
|
||||
opts := retryhelper.DefaultRetryOptions()
|
||||
|
||||
return retryhelper.Retry(func() error {
|
||||
return e.CancelFutOrder(apiUserInfo, symbol, newClientOrderId)
|
||||
}, opts)
|
||||
}
|
||||
|
||||
// CancelFutOrder 通过单个订单号取消合约委托
|
||||
// symbol 交易对
|
||||
// newClientOrderId 系统自定义订单号
|
||||
|
||||
@ -27,7 +27,7 @@ import (
|
||||
/*
|
||||
修改订单信息
|
||||
*/
|
||||
func ChangeFutureOrder(mapData map[string]interface{}) {
|
||||
func ChangeFutureOrder(mapData map[string]interface{}, apiKey string) {
|
||||
// 检查订单号是否存在
|
||||
orderSn, ok := mapData["c"]
|
||||
originOrderSn := mapData["C"] //取消操作 代表原始订单号
|
||||
@ -61,34 +61,45 @@ func ChangeFutureOrder(mapData map[string]interface{}) {
|
||||
}
|
||||
defer lock.Release()
|
||||
|
||||
// 查询订单
|
||||
preOrder, err := getPreOrder(db, orderSn)
|
||||
//反单逻辑
|
||||
reverseService := ReverseService{}
|
||||
reverseService.Orm = db
|
||||
reverseService.Log = logger.NewHelper(logger.DefaultLogger)
|
||||
_, err = reverseService.ReverseOrder(apiKey, mapData)
|
||||
|
||||
if err != nil {
|
||||
logger.Error("合约订单回调失败,查询订单失败:", orderSn, " err:", err)
|
||||
return
|
||||
logger.Errorf("合约订单回调失败,反单失败:%v", err)
|
||||
}
|
||||
|
||||
// 解析订单状态
|
||||
status, ok := mapData["X"].(string)
|
||||
if !ok {
|
||||
mapStr, _ := sonic.Marshal(&mapData)
|
||||
logger.Error("订单回调失败,没有状态:", string(mapStr))
|
||||
return
|
||||
}
|
||||
// 更新订单状态
|
||||
orderStatus, reason := parseOrderStatus(status, mapData)
|
||||
// 以前的下单逻辑
|
||||
// 查询订单
|
||||
// preOrder, err := getPreOrder(db, orderSn)
|
||||
// if err != nil {
|
||||
// logger.Error("合约订单回调失败,查询订单失败:", orderSn, " err:", err)
|
||||
// return
|
||||
// }
|
||||
|
||||
if orderStatus == 0 {
|
||||
logger.Error("订单回调失败,状态错误:", orderSn, " status:", status, " reason:", reason)
|
||||
return
|
||||
}
|
||||
// // 解析订单状态
|
||||
// status, ok := mapData["X"].(string)
|
||||
// if !ok {
|
||||
// mapStr, _ := sonic.Marshal(&mapData)
|
||||
// logger.Error("订单回调失败,没有状态:", string(mapStr))
|
||||
// return
|
||||
// }
|
||||
// // 更新订单状态
|
||||
// orderStatus, reason := parseOrderStatus(status, mapData)
|
||||
|
||||
if err := updateOrderStatus(db, preOrder, orderStatus, reason, true, mapData); err != nil {
|
||||
logger.Error("修改订单状态失败:", orderSn, " err:", err)
|
||||
return
|
||||
}
|
||||
// if orderStatus == 0 {
|
||||
// logger.Error("订单回调失败,状态错误:", orderSn, " status:", status, " reason:", reason)
|
||||
// return
|
||||
// }
|
||||
|
||||
handleFutOrderByType(db, preOrder, orderStatus)
|
||||
// if err := updateOrderStatus(db, preOrder, orderStatus, reason, true, mapData); err != nil {
|
||||
// logger.Error("修改订单状态失败:", orderSn, " err:", err)
|
||||
// return
|
||||
// }
|
||||
|
||||
// handleFutOrderByType(db, preOrder, orderStatus)
|
||||
}
|
||||
|
||||
// 合约回调
|
||||
|
||||
@ -61,16 +61,17 @@ type EntryPriceResult struct {
|
||||
}
|
||||
|
||||
type FutOrderPlace struct {
|
||||
ApiId int `json:"api_id"` //api用户id
|
||||
Symbol string `json:"symbol"` //合约交易对
|
||||
Side string `json:"side"` //购买方向
|
||||
Quantity decimal.Decimal `json:"quantity"` //数量
|
||||
Price decimal.Decimal `json:"price"` //限价单价
|
||||
SideType string `json:"side_type"` //现价或者市价
|
||||
OpenOrder int `json:"open_order"` //是否开启限价单止盈止损
|
||||
Profit decimal.Decimal `json:"profit"` //止盈价格
|
||||
StopPrice decimal.Decimal `json:"stopprice"` //止损价格
|
||||
OrderType string `json:"order_type"` //订单类型,市价或限价MARKET(市价单) TAKE_PROFIT_MARKET(市价止盈) TAKE_PROFIT(限价止盈) STOP (限价止损) STOP_MARKET(市价止损)
|
||||
ApiId int `json:"api_id"` //api用户id
|
||||
Symbol string `json:"symbol"` //合约交易对
|
||||
Side string `json:"side"` //购买方向
|
||||
PositionSide string `json:"position_side"` //持仓方向
|
||||
Quantity decimal.Decimal `json:"quantity"` //数量
|
||||
Price decimal.Decimal `json:"price"` //限价单价
|
||||
SideType string `json:"side_type"` //现价或者市价
|
||||
OpenOrder int `json:"open_order"` //是否开启限价单止盈止损
|
||||
Profit decimal.Decimal `json:"profit"` //止盈价格
|
||||
StopPrice decimal.Decimal `json:"stopprice"` //止损价格
|
||||
OrderType string `json:"order_type"` //订单类型,市价或限价MARKET(市价单) TAKE_PROFIT_MARKET(市价止盈) TAKE_PROFIT(限价止盈) STOP (限价止损) STOP_MARKET(市价止损)
|
||||
NewClientOrderId string `json:"newClientOrderId"`
|
||||
}
|
||||
|
||||
|
||||
623
services/binanceservice/reverse_service.go
Normal file
623
services/binanceservice/reverse_service.go
Normal file
@ -0,0 +1,623 @@
|
||||
package binanceservice
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"errors"
|
||||
"fmt"
|
||||
DbModels "go-admin/app/admin/models"
|
||||
"go-admin/common/global"
|
||||
"go-admin/common/helper"
|
||||
"go-admin/pkg/maphelper"
|
||||
"go-admin/services/cacheservice"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/go-admin-team/go-admin-core/sdk/service"
|
||||
"github.com/shopspring/decimal"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type ReverseService struct {
|
||||
service.Service
|
||||
}
|
||||
|
||||
// SaveMainOrder 保存订单信息
|
||||
// mapData: 下单数据
|
||||
// status: 订单状态
|
||||
// orderSn: 自定义订单号
|
||||
func (e *ReverseService) handleSimpleStatusChange(status string, orderSn string, mapData map[string]interface{}) (bool, error) {
|
||||
statusMap := map[string]int{
|
||||
"NEW": 2,
|
||||
"CANCELED": 6,
|
||||
"EXPIRED": 7,
|
||||
"EXPIRED_IN_MATCH": 7,
|
||||
}
|
||||
if newStatus, ok := statusMap[status]; ok {
|
||||
_ = e.changeOrderStatus(newStatus, orderSn, mapData)
|
||||
return false, nil
|
||||
}
|
||||
return false, fmt.Errorf("不支持的订单状态 %s", status)
|
||||
}
|
||||
|
||||
// ReverseOrder 反向下单
|
||||
// apiKey: api key
|
||||
// mapData: 下单数据
|
||||
// return apiInfo, bool
|
||||
// bool: true=需要反单, false=不需要反单
|
||||
func (e *ReverseService) ReverseOrder(apiKey string, mapData map[string]interface{}) (bool, error) {
|
||||
apiInfo := GetApiInfoByKey(apiKey)
|
||||
|
||||
if apiInfo.Id == 0 {
|
||||
e.Log.Errorf("获取apiInfo失败 %s", apiKey)
|
||||
return false, nil
|
||||
}
|
||||
|
||||
status, _ := maphelper.GetString(mapData, "X")
|
||||
orderSn, err := maphelper.GetString(mapData, "c")
|
||||
if err != nil {
|
||||
return true, err
|
||||
}
|
||||
ot, err := maphelper.GetString(mapData, "ot")
|
||||
|
||||
if err != nil {
|
||||
return true, err
|
||||
}
|
||||
|
||||
switch status {
|
||||
case "NEW", "CANCELED", "EXPIRED", "EXPIRED_IN_MATCH":
|
||||
result, err := e.handleSimpleStatusChange(status, orderSn, mapData)
|
||||
|
||||
//如果是 新开止盈止损 需要取消反单的止盈止损之后重下反单止盈止损
|
||||
if status == "NEW" && (ot == "TAKE_PROFIT_MARKET" || ot == "STOP_LOSS_MARKET" || ot == "TAKE_PROFIT_LIMIT" || ot == "STOP_LOSS_LIMIT") &&
|
||||
apiInfo.ReverseStatus == 1 && apiInfo.OpenStatus == 1 && apiInfo.ReverseApiId > 0 {
|
||||
|
||||
if err := e.ReTakeOrStopOrder(&mapData, orderSn, &apiInfo); err != nil {
|
||||
return true, err
|
||||
}
|
||||
}
|
||||
|
||||
return result, err
|
||||
case "FILLED":
|
||||
if apiInfo.ReverseStatus == 1 && apiInfo.OpenStatus == 1 && apiInfo.ReverseApiId > 0 {
|
||||
reverseApiInfo, err := GetApiInfo(apiInfo.ReverseApiId)
|
||||
if err != nil {
|
||||
e.Log.Errorf("获取反向api信息失败 reverseApiId:%d, err:%v", apiInfo.ReverseApiId, err)
|
||||
return true, err
|
||||
}
|
||||
|
||||
mainOrder, err := e.SaveMainOrder(mapData, apiInfo)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
||||
switch {
|
||||
case mainOrder.PositionSide == "LONG" && mainOrder.Side == "BUY", mainOrder.PositionSide == "SHORT" && mainOrder.Side == "SELL":
|
||||
if mainOrder.Category == 0 {
|
||||
if err1 := e.savePosition(&mainOrder, reverseApiInfo.Id, true, false, false); err1 != nil {
|
||||
return true, err1
|
||||
}
|
||||
|
||||
e.DoAddReverseOrder(&mainOrder, &reverseApiInfo, apiInfo.OrderProportion, false, false)
|
||||
}
|
||||
case mainOrder.PositionSide == "SHORT" && mainOrder.Side == "BUY", mainOrder.PositionSide == "LONG" && mainOrder.Side == "SELL":
|
||||
if mainOrder.Category == 0 {
|
||||
closePosition := maphelper.GetBool(mapData, "R")
|
||||
if err1 := e.savePosition(&mainOrder, reverseApiInfo.Id, true, true, closePosition); err1 != nil {
|
||||
e.Log.Errorf("保存主订单失败: %v", err1)
|
||||
return true, err1
|
||||
}
|
||||
e.DoAddReverseOrder(&mainOrder, &reverseApiInfo, apiInfo.OrderProportion, true, closePosition)
|
||||
}
|
||||
default:
|
||||
return true, errors.New("不支持的订单类型")
|
||||
}
|
||||
return true, nil
|
||||
} else if apiInfo.Subordinate == "2" {
|
||||
symbol, err := maphelper.GetString(mapData, "s")
|
||||
|
||||
if err != nil {
|
||||
return true, err
|
||||
}
|
||||
positionSide, err := maphelper.GetString(mapData, "ps")
|
||||
|
||||
if err != nil {
|
||||
return true, err
|
||||
}
|
||||
side, err := maphelper.GetString(mapData, "S")
|
||||
if err != nil {
|
||||
return true, err
|
||||
}
|
||||
|
||||
totalNum := maphelper.GetDecimal(mapData, "z")
|
||||
|
||||
mainOrder := DbModels.LineReverseOrder{
|
||||
ApiId: apiInfo.Id,
|
||||
Symbol: symbol,
|
||||
PositionSide: positionSide,
|
||||
TotalNum: totalNum,
|
||||
Side: side,
|
||||
}
|
||||
e.changeOrderStatus(3, orderSn, mapData)
|
||||
|
||||
switch {
|
||||
case mainOrder.PositionSide == "LONG" && mainOrder.Side == "BUY", mainOrder.PositionSide == "SHORT" && mainOrder.Side == "SELL":
|
||||
if mainOrder.Category == 0 {
|
||||
if err1 := e.savePosition(&mainOrder, 0, false, false, false); err1 != nil {
|
||||
return true, err1
|
||||
}
|
||||
}
|
||||
case mainOrder.PositionSide == "SHORT" && mainOrder.Side == "BUY", mainOrder.PositionSide == "LONG" && mainOrder.Side == "SELL":
|
||||
if mainOrder.Category == 0 {
|
||||
closePosition := maphelper.GetBool(mapData, "R")
|
||||
if err1 := e.savePosition(&mainOrder, 0, false, true, closePosition); err1 != nil {
|
||||
return true, err1
|
||||
}
|
||||
}
|
||||
default:
|
||||
return true, errors.New("不支持的订单类型")
|
||||
}
|
||||
}
|
||||
default:
|
||||
return false, fmt.Errorf("不支持的订单状态 %s", status)
|
||||
}
|
||||
|
||||
return false, nil
|
||||
}
|
||||
|
||||
// 修改订单状态
|
||||
// status: 订单状态 1-待下单 2-已下单 3-已成交 6-已取消 7-已过期
|
||||
func (e *ReverseService) changeOrderStatus(status int, orderSn string, mapData map[string]interface{}) error {
|
||||
data := map[string]interface{}{"status": status, "updated_at": time.Now()}
|
||||
|
||||
if status == 3 {
|
||||
now := time.Now()
|
||||
if orderId, ok := mapData["i"].(float64); ok {
|
||||
data["order_id"] = orderId
|
||||
}
|
||||
|
||||
if ap, ok := mapData["ap"].(bool); ok {
|
||||
data["final_price"] = ap
|
||||
}
|
||||
|
||||
if num, ok := mapData["z"].(string); ok {
|
||||
data["total_num"], _ = decimal.NewFromString(num)
|
||||
}
|
||||
|
||||
data["trigger_time"] = &now
|
||||
}
|
||||
|
||||
db := e.Orm.Model(&DbModels.LineReverseOrder{}).
|
||||
Where("order_sn =? and status != 3", orderSn).
|
||||
Updates(data)
|
||||
|
||||
if db.Error != nil {
|
||||
e.Log.Errorf("修改订单状态失败 orderSn:%s, err:%v", orderSn, db.Error)
|
||||
}
|
||||
|
||||
if db.RowsAffected == 0 {
|
||||
e.Log.Errorf("修改订单状态失败 orderSn:%s, 未找到订单", orderSn)
|
||||
}
|
||||
return db.Error
|
||||
}
|
||||
|
||||
// 先保存主单,持仓信息 必须用双向持仓!
|
||||
func (e *ReverseService) SaveMainOrder(mapData map[string]interface{}, apiInfo DbModels.LineApiUser) (DbModels.LineReverseOrder, error) {
|
||||
now := time.Now()
|
||||
symbol, err := maphelper.GetString(mapData, "s")
|
||||
var reverseOrder DbModels.LineReverseOrder
|
||||
if err != nil {
|
||||
return reverseOrder, err
|
||||
}
|
||||
|
||||
orderSn, err := maphelper.GetString(mapData, "c")
|
||||
if err != nil {
|
||||
return reverseOrder, err
|
||||
}
|
||||
|
||||
side, err := maphelper.GetString(mapData, "S")
|
||||
if err != nil {
|
||||
return reverseOrder, err
|
||||
}
|
||||
|
||||
positionSide, err := maphelper.GetString(mapData, "ps")
|
||||
if err != nil {
|
||||
return reverseOrder, err
|
||||
}
|
||||
|
||||
mainType, _ := maphelper.GetString(mapData, "ot")
|
||||
|
||||
reverseOrder = DbModels.LineReverseOrder{
|
||||
ApiId: apiInfo.Id,
|
||||
Category: 0,
|
||||
OrderSn: orderSn,
|
||||
TriggerTime: &now,
|
||||
Status: 3,
|
||||
Symbol: symbol,
|
||||
FollowOrderSn: "",
|
||||
Type: mainType,
|
||||
Price: maphelper.GetDecimal(mapData, "p"),
|
||||
FinalPrice: maphelper.GetDecimal(mapData, "ap"),
|
||||
TotalNum: maphelper.GetDecimal(mapData, "z"),
|
||||
Side: side,
|
||||
PositionSide: positionSide,
|
||||
}
|
||||
|
||||
reverseOrder.PriceU = reverseOrder.Price
|
||||
|
||||
if !reverseOrder.Price.IsZero() && !reverseOrder.TotalNum.IsZero() {
|
||||
reverseOrder.BuyPrice = reverseOrder.Price.Mul(reverseOrder.TotalNum)
|
||||
}
|
||||
|
||||
if id, err := maphelper.GetFloat64(mapData, "i"); err == nil {
|
||||
reverseOrder.OrderId = strconv.FormatFloat(id, 'f', -1, 64)
|
||||
}
|
||||
|
||||
orderType, _ := maphelper.GetString(mapData, "ot")
|
||||
switch orderType {
|
||||
case "LIMIT", "MARKET":
|
||||
reverseOrder.OrderType = 0
|
||||
case "TAKE_PROFIT_MARKET", "TAKE_PROFIT":
|
||||
reverseOrder.OrderType = 1
|
||||
case "STOP_MARKET", "STOP", "TRAILING_STOP_MARKET":
|
||||
reverseOrder.OrderType = 2
|
||||
default:
|
||||
return reverseOrder, fmt.Errorf("不支持的订单类型: %s", orderType)
|
||||
}
|
||||
|
||||
if reverseOrder.PositionSide == "BOTH" {
|
||||
return reverseOrder, errors.New("不支持的持仓类型,必须为双向持仓")
|
||||
}
|
||||
|
||||
if err := e.Orm.Create(&reverseOrder).Error; err != nil {
|
||||
e.Log.Errorf("保存主单失败:%v", err)
|
||||
return reverseOrder, err
|
||||
}
|
||||
|
||||
return reverseOrder, nil
|
||||
}
|
||||
|
||||
// 更新仓位信息
|
||||
// closePosition: true=平仓, false=减仓
|
||||
func (e *ReverseService) savePosition(reverseOrder *DbModels.LineReverseOrder, reverseApiId int, isMain, reducePosition, closePosition bool) error {
|
||||
position := DbModels.LineReversePosition{}
|
||||
positionSide := reverseOrder.PositionSide
|
||||
side := reverseOrder.Side
|
||||
totalNum := reverseOrder.TotalNum
|
||||
|
||||
var querySql string
|
||||
sqlStr := ""
|
||||
|
||||
//如果是主单,存储仓位则是反单的持仓方向
|
||||
if isMain {
|
||||
if reverseOrder.PositionSide == "LONG" {
|
||||
positionSide = "SHORT"
|
||||
} else {
|
||||
positionSide = "LONG"
|
||||
}
|
||||
|
||||
if !reducePosition {
|
||||
//反单止盈止损方向相反
|
||||
if side == "SELL" {
|
||||
side = "BUY"
|
||||
} else {
|
||||
side = "SELL"
|
||||
}
|
||||
|
||||
position.ReverseApiId = reverseApiId
|
||||
position.Side = side
|
||||
position.ApiId = reverseOrder.ApiId
|
||||
position.Symbol = reverseOrder.Symbol
|
||||
position.Status = 1
|
||||
position.ReverseStatus = 0
|
||||
position.PositionSide = positionSide
|
||||
}
|
||||
querySql = "api_id =? and position_side =? and symbol =? and status =1"
|
||||
|
||||
//平仓
|
||||
if closePosition {
|
||||
totalNum = decimal.Zero
|
||||
sqlStr = "UPDATE line_reverse_position set amount=@totalNum,updated_at=now(),status=2 where id =@id and status!=2 "
|
||||
} else if reducePosition {
|
||||
//只减仓
|
||||
sqlStr = "UPDATE line_reverse_position set amount=amount - @totalNum,updated_at=now() where id =@id and status!=2 "
|
||||
} else {
|
||||
sqlStr = "UPDATE line_reverse_position set total_amount=total_amount + @totalNum,amount=amount + @totalNum,updated_at=now() where id =@id and status!=2"
|
||||
}
|
||||
} else {
|
||||
querySql = "reverse_api_id =? and position_side =? and symbol =? and reverse_status in (0,1)"
|
||||
|
||||
if closePosition {
|
||||
totalNum = decimal.Zero
|
||||
sqlStr = "UPDATE line_reverse_position set reverse_amount=@totalNum,updated_at=now(),reverse_status=2 where id =@id and reverse_status !=2"
|
||||
} else if reducePosition {
|
||||
sqlStr = "UPDATE line_reverse_position set reverse_amount=reverse_amount - @totalNum,updated_at=now() where id =@id and reverse_status !=2"
|
||||
} else {
|
||||
sqlStr = "UPDATE line_reverse_position set total_reverse_amount=total_reverse_amount + @totalNum,reverse_amount=reverse_amount + @totalNum,updated_at=now(),reverse_status =1 where id =@id and reverse_status !=2"
|
||||
}
|
||||
}
|
||||
|
||||
err := e.Orm.Transaction(func(tx *gorm.DB) error {
|
||||
err1 := tx.Model(&position).Where(querySql,
|
||||
reverseOrder.ApiId, positionSide, reverseOrder.Symbol).First(&position).Error
|
||||
|
||||
if err1 != nil {
|
||||
//主单仓位不存在,创建新仓位
|
||||
if isMain && errors.Is(err1, gorm.ErrRecordNotFound) && !reducePosition {
|
||||
|
||||
if err2 := tx.Create(&position).Error; err2 != nil {
|
||||
return err2
|
||||
}
|
||||
} else {
|
||||
return err1
|
||||
}
|
||||
}
|
||||
|
||||
dbResult := tx.Exec(sqlStr, sql.Named("totalNum", totalNum), sql.Named("id", position.Id))
|
||||
if dbResult.Error != nil {
|
||||
return dbResult.Error
|
||||
}
|
||||
|
||||
if dbResult.RowsAffected == 0 {
|
||||
e.Log.Errorf("减仓数据 是否平仓单:%v :%v", closePosition, reverseOrder)
|
||||
return errors.New("没有找到对应的持仓信息")
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
// 反向下单
|
||||
// mainOrder: 主单信息
|
||||
// reverseApiId: 反向apiId
|
||||
// orderProportion: 反向下单比例
|
||||
// reducePosition: 是否减仓
|
||||
// closePosition: 是否平仓
|
||||
func (e *ReverseService) DoAddReverseOrder(mainOrder *DbModels.LineReverseOrder, reverseApiInfo *DbModels.LineApiUser, orderProportion decimal.Decimal, reducePosition, closePosition bool) error {
|
||||
order := DbModels.LineReverseOrder{}
|
||||
order.ApiId = reverseApiInfo.Id
|
||||
order.Category = 1
|
||||
order.OrderSn = helper.GetOrderNo()
|
||||
order.FollowOrderSn = mainOrder.OrderSn
|
||||
order.Symbol = mainOrder.Symbol
|
||||
order.OrderType = 0
|
||||
|
||||
switch mainOrder.PositionSide {
|
||||
case "LONG":
|
||||
order.PositionSide = "SHORT"
|
||||
case "SHORT":
|
||||
order.PositionSide = "LONG"
|
||||
}
|
||||
|
||||
switch mainOrder.Side {
|
||||
case "SELL":
|
||||
order.Side = "BUY"
|
||||
case "BUY":
|
||||
order.Side = "SELL"
|
||||
default:
|
||||
return fmt.Errorf("不支持的订单类型 side:%s", mainOrder.Side)
|
||||
}
|
||||
|
||||
if reducePosition && closePosition {
|
||||
order.OrderType = 4
|
||||
} else if reducePosition {
|
||||
order.OrderType = 3
|
||||
}
|
||||
|
||||
symbol, err := cacheservice.GetTradeSet(global.EXCHANGE_BINANCE, mainOrder.Symbol, 1)
|
||||
|
||||
if err != nil {
|
||||
e.Log.Errorf("获取交易对信息失败 symbol:%s custom:%s :%v", mainOrder.Symbol, mainOrder.OrderSn, err)
|
||||
return err
|
||||
}
|
||||
|
||||
setting, err := cacheservice.GetReverseSetting(e.Orm)
|
||||
|
||||
if err != nil {
|
||||
e.Log.Errorf("获取反单设置失败 symbol:%s custom:%s :%v", mainOrder.Symbol, mainOrder.OrderSn, err)
|
||||
return err
|
||||
}
|
||||
|
||||
signPrice, _ := decimal.NewFromString(symbol.LastPrice)
|
||||
price := signPrice.Truncate(int32(symbol.PriceDigit))
|
||||
|
||||
if price.Cmp(decimal.Zero) <= 0 {
|
||||
e.Log.Errorf("获取最新价格失败 symbol:%s custom:%s 单价小于0", mainOrder.Symbol, mainOrder.OrderSn)
|
||||
return errors.New("获取最新价格失败")
|
||||
}
|
||||
|
||||
var percent decimal.Decimal
|
||||
switch {
|
||||
case order.PositionSide == "LONG" && order.Side == "BUY", order.PositionSide == "SHORT" && order.Side == "BUY":
|
||||
percent = decimal.NewFromInt(100).Add(setting.ReversePremiumRatio)
|
||||
case order.PositionSide == "SHORT" && order.Side == "SELL", order.PositionSide == "LONG" && order.Side == "SELL":
|
||||
percent = decimal.NewFromInt(100).Sub(setting.ReversePremiumRatio)
|
||||
default:
|
||||
return fmt.Errorf("不支持的订单类型 ps:%s, side:%s", order.PositionSide, order.Side)
|
||||
}
|
||||
|
||||
percent = percent.Div(decimal.NewFromInt(100)).Truncate(4)
|
||||
//计算溢价单价
|
||||
price = price.Mul(percent).Truncate(int32(symbol.PriceDigit))
|
||||
var amount decimal.Decimal
|
||||
|
||||
//平仓单直接卖出全部数量
|
||||
if closePosition {
|
||||
var position DbModels.LineReversePosition
|
||||
if err1 := e.Orm.Model(position).
|
||||
Where("position_side =? and reverse_api_id =? and reverse_status =1", order.PositionSide, order.ApiId).
|
||||
Select("reverse_amount").
|
||||
Find(&position).Error; err1 != nil {
|
||||
e.Log.Errorf("获取剩余仓位失败 symbol:%s custom:%s :%v", order.Symbol, order.OrderSn, err1)
|
||||
|
||||
if len(err1.Error()) < 255 {
|
||||
order.Remark = err1.Error()
|
||||
} else {
|
||||
order.Remark = err1.Error()[:254]
|
||||
}
|
||||
}
|
||||
|
||||
amount = position.ReverseAmount
|
||||
|
||||
if amount.IsZero() {
|
||||
order.Remark = "没有持仓数量"
|
||||
}
|
||||
} else {
|
||||
proportion := decimal.NewFromInt(100)
|
||||
|
||||
if orderProportion.Cmp(decimal.Zero) > 0 {
|
||||
proportion = orderProportion
|
||||
}
|
||||
|
||||
//反向下单百分比
|
||||
proportion = proportion.Div(decimal.NewFromInt(100)).Truncate(2)
|
||||
amount = mainOrder.BuyPrice.Mul(proportion).Div(price).Truncate(int32(symbol.AmountDigit))
|
||||
|
||||
if amount.Cmp(decimal.Zero) <= 0 {
|
||||
e.Log.Errorf("计算数量失败 symbol:%s custom:%s 数量小于0", mainOrder.Symbol, mainOrder.OrderSn)
|
||||
return errors.New("计算数量失败")
|
||||
}
|
||||
}
|
||||
|
||||
order.TotalNum = amount
|
||||
order.Price = price
|
||||
order.PriceU = price
|
||||
order.BuyPrice = amount.Mul(price).Truncate(int32(symbol.PriceDigit))
|
||||
order.Status = 1
|
||||
order.Type = setting.ReverseOrderType
|
||||
order.SignPrice = signPrice
|
||||
|
||||
if order.Remark != "" {
|
||||
order.Status = 8
|
||||
}
|
||||
|
||||
if err := e.Orm.Model(&order).Create(&order).Error; err != nil {
|
||||
e.Log.Errorf("保存反单失败 symbol:%s custom:%s :%v", mainOrder.Symbol, mainOrder.OrderSn, err)
|
||||
return err
|
||||
}
|
||||
|
||||
if order.Status == 1 {
|
||||
e.DoBianceOrder(&order, reverseApiInfo, &setting, reducePosition, closePosition)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// 处理币安订单
|
||||
// order: 反单信息
|
||||
// apiInfo: 币安api信息
|
||||
func (e *ReverseService) DoBianceOrder(order *DbModels.LineReverseOrder, apiInfo *DbModels.LineApiUser, setting *DbModels.LineReverseSetting, reducePosition, closePosition bool) error {
|
||||
futApiV2 := FuturesResetV2{Service: e.Service}
|
||||
orderType := setting.ReverseOrderType
|
||||
|
||||
if orderType == "" {
|
||||
orderType = "LIMIT"
|
||||
}
|
||||
|
||||
params := FutOrderPlace{
|
||||
ApiId: apiInfo.Id,
|
||||
Symbol: order.Symbol,
|
||||
PositionSide: order.PositionSide,
|
||||
Side: order.Side,
|
||||
OrderType: orderType,
|
||||
Quantity: order.TotalNum,
|
||||
Price: order.Price,
|
||||
NewClientOrderId: order.OrderSn,
|
||||
}
|
||||
|
||||
err := futApiV2.OrderPlaceLoop(apiInfo, params)
|
||||
|
||||
if err != nil {
|
||||
e.Log.Errorf("币安下单失败 symbol:%s custom:%s :%v", order.Symbol, order.OrderSn, err)
|
||||
remark := err.Error()
|
||||
|
||||
if len(remark) > 255 {
|
||||
remark = remark[:254]
|
||||
}
|
||||
|
||||
if err1 := e.Orm.Model(&order).Where("id =? and status !=3", order.Id).Updates(map[string]interface{}{"status": 8, "remark": remark, "updated_at": time.Now()}).Error; err1 != nil {
|
||||
e.Log.Errorf("更新订单状态失败 symbol:%s custom:%s :%v", err1)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// 重下止盈止损
|
||||
// mapData: 主单止盈止损回调
|
||||
func (e *ReverseService) ReTakeOrStopOrder(mapData *map[string]interface{}, orderSn string, mainApiInfo *DbModels.LineApiUser) error {
|
||||
symbol, err := maphelper.GetString(*mapData, "s")
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
side, err := maphelper.GetString(*mapData, "S")
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
ot, err := maphelper.GetString(*mapData, "ot")
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
//反单止盈止损方向相反
|
||||
if side == "SELL" {
|
||||
side = "BUY"
|
||||
} else {
|
||||
side = "SELL"
|
||||
}
|
||||
|
||||
positionSide, err := maphelper.GetString(*mapData, "ps")
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if positionSide == "LONG" {
|
||||
positionSide = "SHORT"
|
||||
} else {
|
||||
positionSide = "LONG"
|
||||
}
|
||||
|
||||
apiInfo, err := GetApiInfo(mainApiInfo.ReverseApiId)
|
||||
|
||||
if err != nil {
|
||||
e.Log.Errorf("根据主单api获取反单api失败 symbol:%s custom:%s :%v", symbol, orderSn, err)
|
||||
return err
|
||||
}
|
||||
|
||||
var reverseOrder DbModels.LineReverseOrder
|
||||
e.Orm.Model(&DbModels.LineReverseOrder{}).
|
||||
Where("symbol =? and api_id =? and position_side =? and side = ? and status =1", symbol, mainApiInfo.ReverseApiId, positionSide, side).
|
||||
First(&reverseOrder)
|
||||
|
||||
//取消旧止盈止损
|
||||
if reverseOrder.OrderSn != "" {
|
||||
futApi := FutRestApi{}
|
||||
err := futApi.CancelFutOrderRetry(apiInfo, symbol, reverseOrder.OrderSn)
|
||||
|
||||
if err != nil {
|
||||
e.Log.Errorf("币安撤单失败 symbol:%s custom:%s :%v", symbol, orderSn, err)
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
params := FutOrderPlace{
|
||||
ApiId: apiInfo.Id,
|
||||
Symbol: symbol,
|
||||
PositionSide: positionSide,
|
||||
Side: side,
|
||||
OrderType: ot,
|
||||
Quantity: reverseOrder.TotalNum,
|
||||
Price: reverseOrder.Price,
|
||||
NewClientOrderId: reverseOrder.OrderSn,
|
||||
}
|
||||
futApiV2 := FuturesResetV2{Service: e.Service}
|
||||
futApiV2.OrderPlace(&apiInfo, params)
|
||||
return nil
|
||||
}
|
||||
1
services/binanceservice/reverse_service_test.go
Normal file
1
services/binanceservice/reverse_service_test.go
Normal file
@ -0,0 +1 @@
|
||||
package binanceservice
|
||||
@ -21,7 +21,6 @@ import (
|
||||
|
||||
"github.com/bytedance/sonic"
|
||||
"github.com/go-admin-team/go-admin-core/logger"
|
||||
log "github.com/go-admin-team/go-admin-core/logger"
|
||||
"github.com/shopspring/decimal"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
@ -29,7 +28,7 @@ import (
|
||||
/*
|
||||
订单回调
|
||||
*/
|
||||
func ChangeSpotOrder(mapData map[string]interface{}) {
|
||||
func ChangeSpotOrder(mapData map[string]interface{}, apiKey string) {
|
||||
// 检查订单号是否存在
|
||||
orderSn, ok := mapData["c"]
|
||||
originOrderSn := mapData["C"] //取消操作 代表原始订单号
|
||||
@ -193,7 +192,7 @@ func handleMainReduceFilled(db *gorm.DB, preOrder *DbModels.LinePreOrder) {
|
||||
lock := helper.NewRedisLock(fmt.Sprintf(rediskey.SpotReduceCallback, preOrder.ApiId, preOrder.Symbol), 120, 20, 100*time.Millisecond)
|
||||
|
||||
if ok, err := lock.AcquireWait(context.Background()); err != nil {
|
||||
log.Error("获取锁失败", err)
|
||||
logger.Error("获取锁失败", err)
|
||||
return
|
||||
} else if ok {
|
||||
defer lock.Release()
|
||||
|
||||
23
services/cacheservice/confg_server_test.go
Normal file
23
services/cacheservice/confg_server_test.go
Normal file
@ -0,0 +1,23 @@
|
||||
package cacheservice
|
||||
|
||||
import (
|
||||
"go-admin/common/helper"
|
||||
"testing"
|
||||
|
||||
"gorm.io/driver/mysql"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
func TestGetReverseSetting(t *testing.T) {
|
||||
dsn := "root:123456@tcp(127.0.0.1:3306)/go_exchange_single?charset=utf8mb4&parseTime=True&loc=Local&timeout=1000ms"
|
||||
db, _ := gorm.Open(mysql.Open(dsn), &gorm.Config{})
|
||||
helper.InitDefaultRedis("127.0.0.1:6379", "", 2)
|
||||
|
||||
setting, err := GetReverseSetting(db)
|
||||
|
||||
if err != nil {
|
||||
t.Error(err)
|
||||
}
|
||||
|
||||
t.Log(setting)
|
||||
}
|
||||
@ -130,3 +130,22 @@ func ResetSystemSetting(db *gorm.DB) (models.LineSystemSetting, error) {
|
||||
}
|
||||
return models.LineSystemSetting{}, nil
|
||||
}
|
||||
|
||||
// 获取反单配置
|
||||
func GetReverseSetting(db *gorm.DB) (models.LineReverseSetting, error) {
|
||||
setting := models.LineReverseSetting{}
|
||||
|
||||
helper.DefaultRedis.HGetAsObject(rediskey.ReverseSetting, &setting)
|
||||
|
||||
if setting.ReverseOrderType == "" {
|
||||
if err := db.Model(&setting).First(&setting).Error; err != nil {
|
||||
logger.Error("获取反单配置失败", err)
|
||||
}
|
||||
|
||||
if err := helper.DefaultRedis.SetHashWithTags(rediskey.ReverseSetting, &setting); err != nil {
|
||||
logger.Error("redis添加反单配置失败", err)
|
||||
}
|
||||
}
|
||||
|
||||
return setting, nil
|
||||
}
|
||||
|
||||
@ -15,7 +15,7 @@ import (
|
||||
- @msg 消息内容
|
||||
- @listenType 订阅类型 0-现货 1-合约
|
||||
*/
|
||||
func ReceiveListen(msg []byte, listenType int) (reconnect bool, err error) {
|
||||
func ReceiveListen(msg []byte, listenType int, apiKey string) (reconnect bool, err error) {
|
||||
var dataMap map[string]interface{}
|
||||
err = sonic.Unmarshal(msg, &dataMap)
|
||||
|
||||
@ -52,7 +52,7 @@ func ReceiveListen(msg []byte, listenType int) (reconnect bool, err error) {
|
||||
}
|
||||
|
||||
utility.SafeGo(func() {
|
||||
binanceservice.ChangeSpotOrder(mapData)
|
||||
binanceservice.ChangeSpotOrder(mapData, apiKey)
|
||||
})
|
||||
} else {
|
||||
var data futuresdto.OrderTradeUpdate
|
||||
@ -64,7 +64,7 @@ func ReceiveListen(msg []byte, listenType int) (reconnect bool, err error) {
|
||||
}
|
||||
|
||||
utility.SafeGo(func() {
|
||||
binanceservice.ChangeFutureOrder(data.OrderDetails)
|
||||
binanceservice.ChangeFutureOrder(data.OrderDetails, apiKey)
|
||||
})
|
||||
}
|
||||
//订单更新
|
||||
@ -72,9 +72,9 @@ func ReceiveListen(msg []byte, listenType int) (reconnect bool, err error) {
|
||||
log.Info("executionReport 推送:", string(msg))
|
||||
|
||||
if listenType == 0 { //现货
|
||||
binanceservice.ChangeSpotOrder(dataMap)
|
||||
binanceservice.ChangeSpotOrder(dataMap, apiKey)
|
||||
} else if listenType == 1 { //合约
|
||||
binanceservice.ChangeFutureOrder(dataMap)
|
||||
binanceservice.ChangeFutureOrder(dataMap, apiKey)
|
||||
} else {
|
||||
log.Error("executionReport 不支持的订阅类型", strconv.Itoa(listenType))
|
||||
}
|
||||
@ -101,6 +101,7 @@ func ReceiveListen(msg []byte, listenType int) (reconnect bool, err error) {
|
||||
|
||||
case "eventStreamTerminated":
|
||||
log.Info("账户数据流被终止 type:", getWsTypeName(listenType))
|
||||
return true, nil
|
||||
default:
|
||||
log.Info("未知事件 内容:", string(msg))
|
||||
log.Info("未知事件", event)
|
||||
|
||||
@ -20,6 +20,7 @@ import (
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/bytedance/sonic"
|
||||
@ -42,7 +43,9 @@ type BinanceWebSocketManager struct {
|
||||
isStopped bool // 标记 WebSocket 是否已主动停止
|
||||
mu sync.Mutex // 用于控制并发访问 isStopped
|
||||
cancelFunc context.CancelFunc
|
||||
listenKey string // 新增字段
|
||||
listenKey string // 新增字段
|
||||
reconnecting atomic.Bool // 防止重复重连
|
||||
ConnectTime time.Time // 当前连接建立时间
|
||||
}
|
||||
|
||||
// 已有连接
|
||||
@ -72,6 +75,10 @@ func NewBinanceWebSocketManager(wsType int, apiKey, apiSecret, proxyType, proxyA
|
||||
}
|
||||
}
|
||||
|
||||
func (wm *BinanceWebSocketManager) GetKey() string {
|
||||
return wm.apiKey
|
||||
}
|
||||
|
||||
func (wm *BinanceWebSocketManager) Start() {
|
||||
utility.SafeGo(wm.run)
|
||||
// wm.run()
|
||||
@ -91,14 +98,32 @@ func (wm *BinanceWebSocketManager) Restart(apiKey, apiSecret, proxyType, proxyAd
|
||||
wm.isStopped = false
|
||||
utility.SafeGo(wm.run)
|
||||
} else {
|
||||
wm.reconnect <- struct{}{}
|
||||
log.Warnf("调用restart")
|
||||
wm.triggerReconnect(true)
|
||||
}
|
||||
|
||||
return wm
|
||||
}
|
||||
|
||||
// 触发重连
|
||||
func (wm *BinanceWebSocketManager) triggerReconnect(force bool) {
|
||||
if force {
|
||||
wm.reconnecting.Store(false) // 强制重置标志位
|
||||
}
|
||||
|
||||
if wm.reconnecting.CompareAndSwap(false, true) {
|
||||
log.Warnf("准备重连 key: %s wsType: %v", wm.apiKey, wm.wsType)
|
||||
// 发送信号触发重连协程
|
||||
select {
|
||||
case wm.reconnect <- struct{}{}:
|
||||
default:
|
||||
// 防止阻塞,如果通道满了就跳过
|
||||
}
|
||||
}
|
||||
}
|
||||
func Restart(wm *BinanceWebSocketManager) {
|
||||
wm.reconnect <- struct{}{}
|
||||
log.Warnf("调用restart")
|
||||
wm.triggerReconnect(true)
|
||||
}
|
||||
func (wm *BinanceWebSocketManager) run() {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
@ -201,6 +226,8 @@ func (wm *BinanceWebSocketManager) connect(ctx context.Context) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// 连接成功,更新连接时间
|
||||
wm.ConnectTime = time.Now()
|
||||
log.Info(fmt.Sprintf("已连接到 Binance %s WebSocket【%s】 key:%s", getWsTypeName(wm.wsType), wm.apiKey, listenKey))
|
||||
|
||||
// Ping处理
|
||||
@ -222,10 +249,88 @@ func (wm *BinanceWebSocketManager) connect(ctx context.Context) error {
|
||||
return nil
|
||||
})
|
||||
|
||||
// utility.SafeGoParam(wm.restartConnect, ctx)
|
||||
utility.SafeGo(func() { wm.startListenKeyRenewal2(ctx) })
|
||||
utility.SafeGo(func() { wm.readMessages(ctx) })
|
||||
utility.SafeGo(func() { wm.handleReconnect(ctx) })
|
||||
utility.SafeGo(func() { wm.startPingLoop(ctx) })
|
||||
// utility.SafeGo(func() { wm.startDeadCheck(ctx) })
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// ReplaceConnection 创建新连接并关闭旧连接,实现无缝连接替换
|
||||
func (wm *BinanceWebSocketManager) ReplaceConnection() error {
|
||||
wm.mu.Lock()
|
||||
if wm.isStopped {
|
||||
wm.mu.Unlock()
|
||||
return errors.New("WebSocket 已停止")
|
||||
}
|
||||
oldCtxCancel := wm.cancelFunc
|
||||
oldConn := wm.ws
|
||||
wm.mu.Unlock()
|
||||
|
||||
log.Infof("🔄 正在替换连接: %s", wm.apiKey)
|
||||
|
||||
// 步骤 1:先获取新的 listenKey 和连接
|
||||
newListenKey, err := wm.getListenKey()
|
||||
if err != nil {
|
||||
return fmt.Errorf("获取新 listenKey 失败: %w", err)
|
||||
}
|
||||
|
||||
dialer, err := wm.getDialer()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
newURL := fmt.Sprintf("%s/%s", wm.url, newListenKey)
|
||||
newConn, _, err := dialer.Dial(newURL, nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("新连接 Dial 失败: %w", err)
|
||||
}
|
||||
|
||||
// 步骤 2:创建新的上下文并启动协程
|
||||
newCtx, newCancel := context.WithCancel(context.Background())
|
||||
|
||||
// 设置 ping handler
|
||||
newConn.SetPingHandler(func(appData string) error {
|
||||
log.Infof("收到 Ping(新连接) key:%s msg:%s", wm.apiKey, appData)
|
||||
for x := 0; x < 5; x++ {
|
||||
if err := newConn.WriteControl(websocket.PongMessage, []byte(appData), time.Now().Add(10*time.Second)); err != nil {
|
||||
log.Errorf("Pong 失败 %d 次 err:%v", x, err)
|
||||
time.Sleep(time.Second)
|
||||
continue
|
||||
}
|
||||
break
|
||||
}
|
||||
setLastTime(wm)
|
||||
return nil
|
||||
})
|
||||
|
||||
// 步骤 3:安全切换连接
|
||||
wm.mu.Lock()
|
||||
wm.ws = newConn
|
||||
wm.listenKey = newListenKey
|
||||
wm.ConnectTime = time.Now()
|
||||
wm.cancelFunc = newCancel
|
||||
wm.mu.Unlock()
|
||||
|
||||
log.Infof("✅ 替换连接成功: %s listenKey: %s", wm.apiKey, newListenKey)
|
||||
|
||||
// 步骤 4:启动新连接协程
|
||||
go wm.startListenKeyRenewal2(newCtx)
|
||||
go wm.readMessages(newCtx)
|
||||
go wm.handleReconnect(newCtx)
|
||||
go wm.startPingLoop(newCtx)
|
||||
// go wm.startDeadCheck(newCtx)
|
||||
|
||||
// 步骤 5:关闭旧连接、取消旧协程
|
||||
if oldCtxCancel != nil {
|
||||
oldCtxCancel()
|
||||
}
|
||||
if oldConn != nil {
|
||||
_ = oldConn.Close()
|
||||
log.Infof("🔒 旧连接已关闭: %s", wm.apiKey)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
@ -251,6 +356,7 @@ func setLastTime(wm *BinanceWebSocketManager) {
|
||||
if val != "" {
|
||||
helper.DefaultRedis.SetString(subKey, val)
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
func (wm *BinanceWebSocketManager) getDialer() (*websocket.Dialer, error) {
|
||||
@ -357,7 +463,8 @@ func (wm *BinanceWebSocketManager) readMessages(ctx context.Context) {
|
||||
_, msg, err := wm.ws.ReadMessage()
|
||||
if err != nil && strings.Contains(err.Error(), "websocket: close") {
|
||||
if !wm.isStopped {
|
||||
wm.reconnect <- struct{}{}
|
||||
log.Error("收到关闭消息", err.Error())
|
||||
wm.triggerReconnect(false)
|
||||
}
|
||||
|
||||
log.Error("websocket 关闭")
|
||||
@ -375,8 +482,9 @@ func (wm *BinanceWebSocketManager) readMessages(ctx context.Context) {
|
||||
func (wm *BinanceWebSocketManager) handleOrderUpdate(msg []byte) {
|
||||
setLastTime(wm)
|
||||
|
||||
if reconnect, _ := ReceiveListen(msg, wm.wsType); reconnect {
|
||||
wm.reconnect <- struct{}{}
|
||||
if reconnect, _ := ReceiveListen(msg, wm.wsType, wm.apiKey); reconnect {
|
||||
log.Errorf("收到重连请求")
|
||||
wm.triggerReconnect(false)
|
||||
}
|
||||
}
|
||||
|
||||
@ -447,6 +555,7 @@ func (wm *BinanceWebSocketManager) handleReconnect(ctx context.Context) {
|
||||
retryCount++
|
||||
|
||||
if retryCount >= maxRetries {
|
||||
wm.reconnecting.Store(false)
|
||||
log.Error("重连失败次数过多,退出重连逻辑")
|
||||
return
|
||||
}
|
||||
@ -455,51 +564,88 @@ func (wm *BinanceWebSocketManager) handleReconnect(ctx context.Context) {
|
||||
continue
|
||||
}
|
||||
|
||||
// 重连成功,清除标记
|
||||
wm.reconnecting.Store(false)
|
||||
retryCount = 0
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 定期删除listenkey 并重启ws
|
||||
func (wm *BinanceWebSocketManager) startListenKeyRenewal(ctx context.Context, listenKey string) {
|
||||
time.Sleep(30 * time.Minute)
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
// 假死检测
|
||||
func (wm *BinanceWebSocketManager) DeadCheck() {
|
||||
subKey := fmt.Sprintf(global.USER_SUBSCRIBE, wm.apiKey)
|
||||
val, _ := helper.DefaultRedis.GetString(subKey)
|
||||
if val == "" {
|
||||
log.Warnf("没有订阅信息,无法进行假死检测")
|
||||
return
|
||||
default:
|
||||
if err := wm.deleteListenKey(listenKey); err != nil {
|
||||
log.Error("Failed to renew listenKey: ,type:%v key: %s", wm.wsType, wm.apiKey, err)
|
||||
} else {
|
||||
log.Debug("Successfully delete listenKey")
|
||||
wm.reconnect <- struct{}{}
|
||||
}
|
||||
}
|
||||
|
||||
// ticker := time.NewTicker(5 * time.Minute)
|
||||
// defer ticker.Stop()
|
||||
var data binancedto.UserSubscribeState
|
||||
_ = sonic.Unmarshal([]byte(val), &data)
|
||||
|
||||
// for {
|
||||
// select {
|
||||
// case <-ticker.C:
|
||||
// if wm.isStopped {
|
||||
// return
|
||||
// }
|
||||
var lastTime *time.Time
|
||||
if wm.wsType == 0 {
|
||||
lastTime = data.SpotLastTime
|
||||
} else {
|
||||
lastTime = data.FuturesLastTime
|
||||
}
|
||||
|
||||
// if err := wm.deleteListenKey(listenKey); err != nil {
|
||||
// log.Error("Failed to renew listenKey: ,type:%v key: %s", wm.wsType, wm.apiKey, err)
|
||||
// } else {
|
||||
// log.Debug("Successfully delete listenKey")
|
||||
// wm.reconnect <- struct{}{}
|
||||
// return
|
||||
// }
|
||||
// case <-ctx.Done():
|
||||
// return
|
||||
// }
|
||||
// }
|
||||
// 定义最大静默时间(超出视为假死)
|
||||
var timeout time.Duration
|
||||
if wm.wsType == 0 {
|
||||
timeout = 40 * time.Second // Spot 每 20s ping,40s 足够
|
||||
} else {
|
||||
timeout = 6 * time.Minute // Futures 每 3 分钟 ping
|
||||
}
|
||||
|
||||
if lastTime != nil && time.Since(*lastTime) > timeout {
|
||||
log.Warnf("检测到假死连接 key:%s type:%v, 距离上次通信: %v, 触发重连", wm.apiKey, wm.wsType, time.Since(*lastTime))
|
||||
wm.triggerReconnect(true)
|
||||
}
|
||||
}
|
||||
|
||||
// 主动心跳发送机制
|
||||
func (wm *BinanceWebSocketManager) startPingLoop(ctx context.Context) {
|
||||
ticker := time.NewTicker(1 * time.Minute)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
if wm.isStopped {
|
||||
return
|
||||
}
|
||||
err := wm.ws.WriteMessage(websocket.PingMessage, []byte("ping"))
|
||||
if err != nil {
|
||||
log.Error("主动 Ping Binance 失败:", err)
|
||||
wm.triggerReconnect(false)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 定期删除listenkey 并重启ws
|
||||
// func (wm *BinanceWebSocketManager) startListenKeyRenewal(ctx context.Context, listenKey string) {
|
||||
// time.Sleep(30 * time.Minute)
|
||||
|
||||
// select {
|
||||
// case <-ctx.Done():
|
||||
// return
|
||||
// default:
|
||||
// if err := wm.deleteListenKey(listenKey); err != nil {
|
||||
// log.Error("Failed to renew listenKey: ,type:%v key: %s", wm.wsType, wm.apiKey, err)
|
||||
// } else {
|
||||
// log.Debug("Successfully delete listenKey")
|
||||
// wm.triggerReconnect()
|
||||
// }
|
||||
// }
|
||||
|
||||
// }
|
||||
|
||||
// 定时续期
|
||||
func (wm *BinanceWebSocketManager) startListenKeyRenewal2(ctx context.Context) {
|
||||
ticker := time.NewTicker(30 * time.Minute)
|
||||
@ -527,49 +673,34 @@ func (wm *BinanceWebSocketManager) startListenKeyRenewal2(ctx context.Context) {
|
||||
/*
|
||||
删除listenkey
|
||||
*/
|
||||
func (wm *BinanceWebSocketManager) deleteListenKey(listenKey string) error {
|
||||
client, err := wm.createBinanceClient()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// func (wm *BinanceWebSocketManager) deleteListenKey(listenKey string) error {
|
||||
// client, err := wm.createBinanceClient()
|
||||
// if err != nil {
|
||||
// return err
|
||||
// }
|
||||
|
||||
var resp []byte
|
||||
// var resp []byte
|
||||
|
||||
switch wm.wsType {
|
||||
case 0:
|
||||
path := fmt.Sprintf("/api/v3/userDataStream")
|
||||
params := map[string]interface{}{
|
||||
"listenKey": listenKey,
|
||||
}
|
||||
resp, _, err = client.SendSpotRequestByKey(path, "DELETE", params)
|
||||
// switch wm.wsType {
|
||||
// case 0:
|
||||
// path := fmt.Sprintf("/api/v3/userDataStream")
|
||||
// params := map[string]interface{}{
|
||||
// "listenKey": listenKey,
|
||||
// }
|
||||
// resp, _, err = client.SendSpotRequestByKey(path, "DELETE", params)
|
||||
|
||||
log.Debug(fmt.Sprintf("deleteListenKey resp: %s", string(resp)))
|
||||
case 1:
|
||||
resp, _, err = client.SendFuturesRequestByKey("/fapi/v1/listenKey", "DELETE", nil)
|
||||
log.Debug(fmt.Sprintf("deleteListenKey resp: %s", string(resp)))
|
||||
default:
|
||||
return errors.New("unknown ws type")
|
||||
}
|
||||
// log.Debug(fmt.Sprintf("deleteListenKey resp: %s", string(resp)))
|
||||
// case 1:
|
||||
// resp, _, err = client.SendFuturesRequestByKey("/fapi/v1/listenKey", "DELETE", nil)
|
||||
// log.Debug(fmt.Sprintf("deleteListenKey resp: %s", string(resp)))
|
||||
// default:
|
||||
// return errors.New("unknown ws type")
|
||||
// }
|
||||
|
||||
return err
|
||||
}
|
||||
// return err
|
||||
// }
|
||||
|
||||
func (wm *BinanceWebSocketManager) renewListenKey(listenKey string) error {
|
||||
// payloadParam := map[string]interface{}{
|
||||
// "listenKey": listenKey,
|
||||
// "apiKey": wm.apiKey,
|
||||
// }
|
||||
|
||||
// params := map[string]interface{}{
|
||||
// "id": getUUID(),
|
||||
// "method": "userDataStream.ping",
|
||||
// "params": payloadParam,
|
||||
// }
|
||||
|
||||
// if err := wm.ws.WriteJSON(params); err != nil {
|
||||
// return err
|
||||
// }
|
||||
// wm.ws.WriteJSON()
|
||||
client, err := wm.createBinanceClient()
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
Reference in New Issue
Block a user