package trading import ( "context" "io" "sig-pub/api/pb" "sig-pub/pkg/client" "sig-pub/pkg/types" "sig-pub/pkg/zlog" "sync" "time" "google.golang.org/grpc" ) type TradingService struct { marketClientAside *client.TradeInstanceAside exchangeClient pb.ExchangeServiceClient klineStore *KlineStore subKlineLock sync.Mutex subKlineInsts [3][]string subKlineStream grpc.BidiStreamingClient[pb.ReqStreamSubscribeKline, pb.RspStreamSubscribeKline] } func NewTradingService( marketClientAside *client.TradeInstanceAside, exchangeClient pb.ExchangeServiceClient, ) *TradingService { return &TradingService{ marketClientAside: marketClientAside, exchangeClient: exchangeClient, } } // 初始化历史k线, 订阅实时k线 func (svr *TradingService) Init() (err error) { svr.klineStore = NewKlineSeriesStore() // get instance exchangeTradeInsts, err := svr.marketClientAside.ListExchangeTradeInstance(context.Background(), pb.ExchangeType_OKX) if err != nil { return } for _, exInst := range exchangeTradeInsts { svr.subKlineInsts[exInst.Exchange] = append(svr.subKlineInsts[exInst.Exchange], exInst.InstId) } // 订阅k线 go svr.subscribeKlines(false) return } func (svr *TradingService) subscribeKlines(reconnect bool) { defer func() { svr.subKlineStream = nil go svr.subscribeKlines(true) }() if reconnect { zlog.Infof("subscribeKlines will reconnect after 5s") time.Sleep(5 * time.Second) } // svr.subKlineLock.Lock() // defer svr.subKlineLock.Unlock() stream, err := svr.exchangeClient.SubscribeKline(context.Background()) if err != nil { zlog.Error("subscribeKlines reqeust error: ", err) return } if svr.subKlineStream != nil { svr.subKlineStream.CloseSend() } // 发送订阅消息 exchanges := []pb.ExchangeType{pb.ExchangeType_OKX} for _, exchange := range exchanges { instIds := svr.subKlineInsts[exchange] msg := &pb.ReqStreamSubscribeKline{ SubType: pb.SubscribeType_Subscribe, Exchanges: []pb.ExchangeType{exchange}, InstIds: instIds, Intervals: SubscribeKlineIntervals, OnlyConfirm: true, } zlog.Debugf("send stream subscribe kline msg: %#v", msg) if err = stream.Send(msg); err != nil { zlog.Errorf("send stream subscribe kline msg error: %v", msg, err) return } } // 接收消息的goroutine for { msg, err := stream.Recv() if err == io.EOF { zlog.Debugf("subscribeKlines connection server closeed") return } if err != nil { zlog.Error("subscribeKlines recv error: ", err) return } for _, k := range msg.Kline.Klines { kline := new(types.Kline) kline.ParsePBKline(msg.Kline.Exchange, k) zlog.Debugf("recv: streamId=%d, %v, %s, %#v", msg.Kline.StreamId, msg.Kline.Exchange, msg.Kline.InstId, kline) svr.klineStore.Update(msg.Kline.Exchange, msg.Kline.InstId, kline) } } } var SubscribeKlineIntervals = []string{ // string(types.Interval1s), string(types.Interval1m), string(types.Interval3m), string(types.Interval5m), string(types.Interval15m), string(types.Interval30m), string(types.Interval1h), string(types.Interval2h), string(types.Interval4h), string(types.Interval6h), string(types.Interval12h), string(types.Interval1d), string(types.Interval1d), string(types.Interval2d), string(types.Interval3d), string(types.Interval5d), string(types.Interval1w), }