diff --git a/config/exchange.toml b/config/exchange.toml index 4d567a5..c0a42b7 100644 --- a/config/exchange.toml +++ b/config/exchange.toml @@ -18,8 +18,8 @@ receiveBuffer = 4096 marketSubscribeLimit = 16 consumeBatch = 1024 consumeLater = 2000 # 时间到达later或者数据累计到batch触发consume -# httpProxy = "http://192.168.1.5:7890" -httpProxy = "http://10.255.183.209:7890" +httpProxy = "http://192.168.1.5:7890" +# httpProxy = "http://10.255.183.209:7890" # 模拟盘API交易地址如下: # REST:https://www.okx.com diff --git a/internal/exchange/exchange_service.go b/internal/exchange/exchange_service.go index 48c90e6..75f735f 100644 --- a/internal/exchange/exchange_service.go +++ b/internal/exchange/exchange_service.go @@ -539,6 +539,7 @@ func (svc *ExchangeService) HistoryKline(ctx context.Context, req *pb.ReqHistory return } + // todo 交易产品初始化完成检查 exchangeInst, ok := exchange.ExchangeInsts.Load(exchangeInstId) if !ok { err = fmt.Errorf("trade instance not support for exchange: %s for %s", req.InstId, req.Exchange) @@ -547,13 +548,23 @@ func (svc *ExchangeService) HistoryKline(ctx context.Context, req *pb.ReqHistory // k线长度检查 afterTs, beforeTs, count := int64(req.After), int64(req.Before), int64(req.Count) + if count == 0 { + count = 100 + } nowTime := time.Now().UnixMilli() if afterTs == 0 && beforeTs == 0 { // 拉取最新的 - afterTs = nowTime - } - if count == 0 { - count = 100 + liveK := exchangeInst.LiveKline.Get(interval) + lastTs := liveK.Ts + if !liveK.Confirm { + lastTs = intervalAdder(liveK.Ts, -1) + } + if !req.Asc { + afterTs = lastTs + } else { + // 从低到高拉取 + beforeTs = max(intervalAdder(lastTs, -count+1), KlineBefore0) + } } if afterTs == 0 { afterTs = min(intervalAdder(beforeTs, count), nowTime) diff --git a/internal/trading/kline_store.go b/internal/trading/kline_store.go index b404723..36c16bf 100644 --- a/internal/trading/kline_store.go +++ b/internal/trading/kline_store.go @@ -3,6 +3,7 @@ package trading import ( "context" "fmt" + "io" "math" "sig-pub/api/pb" "sig-pub/pkg/data" @@ -19,61 +20,172 @@ import ( ) type KlineStore struct { - exchangeClient pb.ExchangeServiceClient - subscribeKlineIntervals []string - store [3]*collect.ConcurrentMap[string, *KlineStoreInstance] // K线列表: []exchange + exchangeClient pb.ExchangeServiceClient + + store *types.ExchangeState[*collect.ConcurrentMap[string, *KlineStoreInstance]] // K线列表: []exchange + subKlineIntervals []string // 订阅的k线的周期列表 + subKlineInsts *types.ExchangeState[*collect.SyncMap[string, bool]] // 订阅k线中的交易产品列表 + subKlineStream grpc.BidiStreamingClient[pb.ReqStreamSubscribeKline, pb.RspStreamSubscribeKline] // 订阅k线的stream } func NewKlineSeriesStore(exchangeClient pb.ExchangeServiceClient) (kss *KlineStore) { - // 订阅实时k线周期列表 - subscribeKlineIntervals := collect.Map2Slice(types.SupportedIntervals, func(interval types.Interval, _ types.IntervalAdder) string { + kss = &KlineStore{exchangeClient: exchangeClient} + + // 周期列表 + kss.subKlineIntervals = collect.Map2Slice(types.SupportedIntervals, func(interval types.Interval, _ types.IntervalAdder) string { return string(interval) }) - kss = &KlineStore{ - subscribeKlineIntervals: subscribeKlineIntervals, - exchangeClient: exchangeClient, - } - kss.store[pb.ExchangeType_OKX] = collect.NewConcurrentMap[string, *KlineStoreInstance](64, func(s string) string { return s }) - // kss.klines[pb.ExchangeType_BINANCE] = + // 产品列表 + kss.subKlineInsts = types.NewExchangeState0(func() *collect.SyncMap[string, bool] { + return collect.NewSyncMap[string, bool]() + }) + // 各交易所 store 初始化 + kss.store = types.NewExchangeState0(func() *collect.ConcurrentMap[string, *KlineStoreInstance] { + return collect.NewConcurrentMap[string, *KlineStoreInstance](64, func(s string) string { + return s + }) + }) return } func (s *KlineStore) Init() (err error) { - // 拉取已初始化完成交易产品, 初始化k线, 开始订阅k线 + // 连接 exchange kline stream + go s.connectSubscribeKline(false) // 订阅交易产品初始化完成事件 mq.NatsCreateConsumer("trading", mq.StreamExchange, mq.TopicExchangeTradeInstanceInited, func() *mq.PublishExchangeTradeInstanceInited { return new(mq.PublishExchangeTradeInstanceInited) }, func(msg *mq.PublishExchangeTradeInstanceInited) (err error) { // 初始化k线, 开始订阅k线 zlog.Infof("subscribed TopicExchangeTradeInstanceInited: %#v", msg) - go s.subscribeKlines(msg.Exchange, msg.InstId) + go s.initKlineSeries(msg.Exchange, msg.InstId) return }) + + // todo 拉取已初始化完成交易产品, 初始化k线, 开始订阅k线 return } -func (s *KlineStore) subscribeKlines(exchange pb.ExchangeType, instId string) { - storeInst := s.store[exchange].ComputeIfAbsent(instId, func(k string) *KlineStoreInstance { +// connectSubscribeKline 连接exchange订阅实时k线 +func (s *KlineStore) connectSubscribeKline(reconnect bool) { + defer func() { + if s.subKlineStream != nil { + s.subKlineStream.CloseSend() + s.subKlineStream = nil + } + go s.connectSubscribeKline(true) + }() + + if reconnect { + zlog.Infof("subscribeKlines will reconnect after 5s") + time.Sleep(5 * time.Second) + } + + stream, err := s.exchangeClient.SubscribeKline(context.Background()) + if err != nil { + zlog.Error("subscribeKlines reqeust error: ", err) + return + } + s.subKlineStream = stream + + // 发送所有交易产品订阅消息 + go func() { + s.subKlineInsts.Range(func(exchange pb.ExchangeType, m *collect.SyncMap[string, bool]) { + // todo 分批订阅 + var instIds []string + m.Range(func(instId string, _ bool) bool { + instIds = append(instIds, instId) + return true + }) + s.sendSubscribeKline(false, exchange, instIds...) + }) + }() + + // 接收订阅k线消息 + 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) + // kline klineStore -> klineSeries -> strategy -> indicator -> klineSeries.Series + s.Update(msg.Kline.Exchange, msg.Kline.InstId, kline) + } + } +} + +// subscribeKline 发送订阅消息 +func (s *KlineStore) sendSubscribeKline(save bool, exchange pb.ExchangeType, instIds ...string) { + if len(instIds) == 0 { + return + } + // 交易产品订阅记录 + if save { + for _, instId := range instIds { + s.subKlineInsts.Get(exchange).Store(instId, true) + } + } + + // 发送订阅消息 + subMsg := &pb.ReqStreamSubscribeKline{ + SubType: pb.SubscribeType_Subscribe, + Exchanges: []pb.ExchangeType{exchange}, + InstIds: instIds, + Intervals: s.subKlineIntervals, + OnlyConfirm: true, + } + doSend := func(retry uint32) (_ int, err error) { + if s.subKlineStream == nil { + return + } + zlog.Debugf("send stream subscribe kline msg: retry=%d, %#v", retry, subMsg) + if err = s.subKlineStream.Send(subMsg); err != nil { + zlog.Errorf("send stream subscribe kline msg error: %v", subMsg, err) + return + } + return + } + if _, err := doSend(0); err == nil { + return + } + + go retry.DoWithFixDelay(math.MaxInt32, time.Second, doSend) +} + +func (s *KlineStore) initKlineSeries(exchange pb.ExchangeType, instId string) { + if !s.store.IsSupport(exchange) { + return + } + + storeInst := s.store.Get(exchange).ComputeIfAbsent(instId, func(k string) *KlineStoreInstance { return NewKlineStoreInstance(exchange, k) }) // 初始化最新的 klineSeries - for _, interval := range s.subscribeKlineIntervals { + for _, interval := range s.subKlineIntervals { for { - after, before := int64(0), int64(0) + before := int64(0) rsp, err := retry.DoWithFixDelay(math.MaxInt32, 2*time.Second, func(retryTimes uint32) (rsp *pb.RspHistoryKline, err error) { rsp, err = s.exchangeClient.HistoryKline(context.Background(), &pb.ReqHistoryKline{ Exchange: exchange, InstId: instId, Interval: string(interval), Count: MaxSeriesKlines, - After: after, + After: 0, Before: before, Live: false, Asc: true, }, grpc.UseCompressor("snappy")) if err != nil { - zlog.Errorf("fetch missing klines error: instId=%s(%s) after=%d before=%d, err=%v", instId, exchange, after, before, err) + zlog.Errorf("fetch missing klines error: instId=%s(%s) before=%d, err=%v", instId, exchange, before, err) } return }) @@ -81,7 +193,11 @@ func (s *KlineStore) subscribeKlines(exchange pb.ExchangeType, instId string) { zlog.Errorf("trade instance initial failed: %s(%s), %v", instId, exchange, err) return } - for _, kline := range rsp.Klines { + klines := rsp.Klines + if len(klines) == 0 { + break + } + for _, kline := range klines { k := new(types.Kline) k.ParsePBKline(exchange, kline) _, _, err = storeInst.Update(k) @@ -91,24 +207,28 @@ func (s *KlineStore) subscribeKlines(exchange pb.ExchangeType, instId string) { } s.Update(exchange, instId, k) - after = k.Ts } + before = klines[len(klines)-1].Ts + if !rsp.Next { + break + } } } - - // 开始订阅k线 + // 初始化历史k线完成, 开始订阅k线 + storeInst.Status.Store(int32(data.StatusOk)) + s.sendSubscribeKline(true, exchange, instId) } // Update // kline klineStore -> klineSeries -> strategy -> indicator -> klineSeries.Series func (s *KlineStore) Update(exchange pb.ExchangeType, instId string, kline *types.Kline) { - storeInst := s.store[exchange].ComputeIfAbsent(instId, func(k string) *KlineStoreInstance { + storeInst := s.store.Get(exchange).ComputeIfAbsent(instId, func(k string) *KlineStoreInstance { return NewKlineStoreInstance(exchange, k) }) // 只处理已初始化完成的交易产品k线 - if storeInst.status.Load() != int32(data.StatusOk) { + if storeInst.Status.Load() != int32(data.StatusOk) { return } before, serial, err := storeInst.Update(kline) @@ -162,7 +282,7 @@ type KlineStoreInstance struct { Exchange pb.ExchangeType InstId string intervalKlines *types.IntervalState[*KlineSeries] - status atomic.Int32 // 交易产品状态 + Status atomic.Int32 // 交易产品状态 } func NewKlineStoreInstance(exchange pb.ExchangeType, instId string) *KlineStoreInstance { @@ -171,7 +291,7 @@ func NewKlineStoreInstance(exchange pb.ExchangeType, instId string) *KlineStoreI InstId: instId, intervalKlines: types.NewIntervalState[*KlineSeries](), } - si.status.Store(int32(data.StatusProcessing)) + si.Status.Store(int32(data.StatusProcessing)) for interval := range types.SupportedIntervals { si.intervalKlines.Set(interval, NewKlineSeries(exchange, instId, interval)) diff --git a/internal/trading/trading_service.go b/internal/trading/trading_service.go index 57ae2d4..79e8430 100644 --- a/internal/trading/trading_service.go +++ b/internal/trading/trading_service.go @@ -1,16 +1,8 @@ package trading import ( - "context" - "io" "sig-pub/api/pb" "sig-pub/pkg/client" - "sig-pub/pkg/types" - "sig-pub/pkg/utils/collect" - "sig-pub/pkg/zlog" - "time" - - "google.golang.org/grpc" ) type TradingService struct { @@ -18,25 +10,16 @@ type TradingService struct { exchangeClient pb.ExchangeServiceClient klineStore *KlineStore - - subscribeKlineIntervals []string - subKlineInsts [3][]string - subKlineStream grpc.BidiStreamingClient[pb.ReqStreamSubscribeKline, pb.RspStreamSubscribeKline] } func NewTradingService( marketClientAside *client.TradeInstanceAside, exchangeClient pb.ExchangeServiceClient, ) *TradingService { - // 订阅实时k线周期列表 - subscribeKlineIntervals := collect.Map2Slice(types.SupportedIntervals, func(interval types.Interval, _ types.IntervalAdder) string { - return string(interval) - }) return &TradingService{ - marketClientAside: marketClientAside, - exchangeClient: exchangeClient, - klineStore: NewKlineSeriesStore(exchangeClient), - subscribeKlineIntervals: subscribeKlineIntervals, + marketClientAside: marketClientAside, + exchangeClient: exchangeClient, + klineStore: NewKlineSeriesStore(exchangeClient), } } @@ -45,81 +28,5 @@ func (svr *TradingService) Init() (err error) { if err = svr.klineStore.Init(); err != nil { return } - - // 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.subscribeStreamKlines(false) return } - -func (svr *TradingService) subscribeStreamKlines(reconnect bool) { - defer func() { - if svr.subKlineStream != nil { - svr.subKlineStream.CloseSend() - svr.subKlineStream = nil - } - go svr.subscribeStreamKlines(true) - }() - - if reconnect { - zlog.Infof("subscribeKlines will reconnect after 5s") - time.Sleep(5 * time.Second) - } - - stream, err := svr.exchangeClient.SubscribeKline(context.Background()) - if err != nil { - zlog.Error("subscribeKlines reqeust error: ", err) - return - } - svr.subKlineStream = stream - - // 发送订阅消息 - exchanges := []pb.ExchangeType{pb.ExchangeType_OKX} - for _, exchange := range exchanges { - instIds := svr.subKlineInsts[exchange] - if len(instIds) == 0 { - continue - } - msg := &pb.ReqStreamSubscribeKline{ - SubType: pb.SubscribeType_Subscribe, - Exchanges: []pb.ExchangeType{exchange}, - InstIds: instIds, - Intervals: svr.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) - // kline klineStore -> klineSeries -> strategy -> indicator -> klineSeries.Series - svr.klineStore.Update(msg.Kline.Exchange, msg.Kline.InstId, kline) - } - } -} diff --git a/pkg/mq/nats_topic.go b/pkg/mq/nats_topic.go index bd86b00..ae7da68 100644 --- a/pkg/mq/nats_topic.go +++ b/pkg/mq/nats_topic.go @@ -1,6 +1,6 @@ package mq -// 初始化jetStream流, {streamName, topics...} +// 初始化jetStream流, {StreamExchange, topics...} var initialStream = [][]string{ {StreamExchange, "exchange.>"}, } diff --git a/pkg/types/exchange.go b/pkg/types/exchange.go new file mode 100644 index 0000000..d5df50a --- /dev/null +++ b/pkg/types/exchange.go @@ -0,0 +1,56 @@ +package types + +import ( + "sig-pub/api/pb" + "sig-pub/pkg/utils/collect" +) + +var SupportedExchanges = []pb.ExchangeType{ + pb.ExchangeType_OKX, +} + +func IsSupportExchange(exchange pb.ExchangeType) bool { + for _, se := range SupportedExchanges { + if se == exchange { + return true + } + } + return false +} + +type ExchangeState[T any] struct { + state []T +} + +func NewExchangeState[T any]() *ExchangeState[T] { + return NewExchangeState0[T](func() (v T) { return }) +} + +func NewExchangeState0[T any](newer func() T) *ExchangeState[T] { + maxExchange := collect.MustMax(SupportedExchanges, func(e pb.ExchangeType) int32 { return int32(e) }) + es := &ExchangeState[T]{ + state: make([]T, maxExchange+1), + } + for _, exchange := range SupportedExchanges { + es.Set(exchange, newer()) + } + return es +} + +func (s *ExchangeState[T]) IsSupport(exchange pb.ExchangeType) bool { + return IsSupportExchange(exchange) +} + +func (s *ExchangeState[T]) Get(exchange pb.ExchangeType) T { + return s.state[exchange] +} + +func (s *ExchangeState[T]) Set(exchange pb.ExchangeType, v T) { + s.state[exchange] = v +} + +func (s *ExchangeState[T]) Range(f func(exchange pb.ExchangeType, m T)) { + for _, exchange := range SupportedExchanges { + f(exchange, s.state[exchange]) + } +}