From 74586b64acb7072c6f0ed2dea7355ed1424615a1 Mon Sep 17 00:00:00 2001 From: strange Date: Tue, 28 Oct 2025 18:37:24 +0800 Subject: [PATCH] reactor trading --- README.md | 4 +- config/exchange.toml | 4 +- internal/trading/backtest/backtest.go | 21 ++++ internal/trading/backtrace/backtrace.go | 16 ---- .../{kline_store.go => kline_series_store.go} | 95 ++++++++++++------- internal/trading/libs/account_position.go | 3 - internal/trading/sig/account_position.go | 3 + .../trading/{ => sig}/indicator_context.go | 2 +- .../trading/{ => sig}/indicator_series.go | 2 +- .../trading/{ => sig}/indicator_service.go | 2 +- internal/trading/{ => sig}/kline_series.go | 2 +- .../trading/{ => sig}/strategy_context.go | 63 ++++++------ .../trading/{ => sig}/trading_plan_runner.go | 2 +- internal/trading/trading_service.go | 34 +++---- pkg/grpc/client/direct_client_factory.go | 7 +- pkg/strategy/gold_x.go | 8 ++ pkg/strategy/sig_strategy.go | 14 --- pkg/strategy/sig_strategy_exchanges.go | 20 ++++ pkg/strategy/sig_strategy_intervals.go | 19 ++-- pkg/strategy/sig_strategy_registry.go | 6 +- pkg/strategy/strategy.go | 30 ++++++ pkg/utils/lang/condition.go | 19 ++++ 22 files changed, 235 insertions(+), 141 deletions(-) create mode 100644 internal/trading/backtest/backtest.go delete mode 100644 internal/trading/backtrace/backtrace.go rename internal/trading/{kline_store.go => kline_series_store.go} (75%) delete mode 100644 internal/trading/libs/account_position.go create mode 100644 internal/trading/sig/account_position.go rename internal/trading/{ => sig}/indicator_context.go (99%) rename internal/trading/{ => sig}/indicator_series.go (98%) rename internal/trading/{ => sig}/indicator_service.go (98%) rename internal/trading/{ => sig}/kline_series.go (99%) rename internal/trading/{ => sig}/strategy_context.go (56%) rename internal/trading/{ => sig}/trading_plan_runner.go (98%) create mode 100644 pkg/utils/lang/condition.go diff --git a/README.md b/README.md index 3e87630..26abae2 100644 --- a/README.md +++ b/README.md @@ -79,4 +79,6 @@ strategy0: 趋势追踪,增长趋势, ### 量化框架参考 -[investing-algorithm-framework](https://github.com/coding-kitties/investing-algorithm-framework) \ No newline at end of file +[investing-algorithm-framework](https://github.com/coding-kitties/investing-algorithm-framework) + +回测信号可视化, /trading/strategySeries 一样从postgres拉信号/订单数据 diff --git a/config/exchange.toml b/config/exchange.toml index c0a42b7..4d567a5 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/trading/backtest/backtest.go b/internal/trading/backtest/backtest.go new file mode 100644 index 0000000..6d36611 --- /dev/null +++ b/internal/trading/backtest/backtest.go @@ -0,0 +1,21 @@ +package backtest + +import ( + "sig-pub/pkg/data/entity" + "sig-pub/pkg/strategy" +) + +// 回测引擎 +// sig strategy +// trade strategy +// close strategy +// 历史k线加载 +type BacktestEngine struct { + start, end int64 + plan entity.TradePlan // 交易计划 +} + +// 多周期策略回测引擎 +type MultiIntervalBacktraceEngine struct { + strategy strategy.IIntervalSigStrategy +} diff --git a/internal/trading/backtrace/backtrace.go b/internal/trading/backtrace/backtrace.go deleted file mode 100644 index aad92ed..0000000 --- a/internal/trading/backtrace/backtrace.go +++ /dev/null @@ -1,16 +0,0 @@ -package backtrace - -import "sig-pub/pkg/strategy" - -// 回测引擎 -// sig strategy -// trade strategy -// close strategy -type BacktraceEngine struct { - strategy strategy.ISigStrategy -} - -// 多周期策略回测引擎 -type MultiIntervalBacktraceEngine struct { - strategy strategy.IIntervalsSigStrategy -} diff --git a/internal/trading/kline_store.go b/internal/trading/kline_series_store.go similarity index 75% rename from internal/trading/kline_store.go rename to internal/trading/kline_series_store.go index b3e6454..8cd9d77 100644 --- a/internal/trading/kline_store.go +++ b/internal/trading/kline_series_store.go @@ -6,6 +6,7 @@ import ( "io" "math" "sig-pub/api/pb" + "sig-pub/internal/trading/sig" "sig-pub/pkg/data" "sig-pub/pkg/mq" "sig-pub/pkg/strategy" @@ -18,20 +19,20 @@ import ( "google.golang.org/grpc" ) -type KlineStore struct { +type KlineSeriesStore struct { exchangeClient pb.ExchangeServiceClient - store *types.ExchangeState[*collect.ConcurrentMap[string, *TradeInstanceKlineSeries]] // K线列表: []exchange - subKlineIntervals []string // 订阅的k线的周期列表 - subKlineInsts *types.ExchangeState[*collect.SyncMap[string, bool]] // 订阅k线中的交易产品列表 - subKlineStream grpc.BidiStreamingClient[pb.ReqStreamSubscribeKline, pb.RspStreamSubscribeKline] // 订阅k线的stream + store *types.ExchangeState[*collect.ConcurrentMap[string, *sig.TradeInstanceKlineSeries]] // K线列表: []exchange + subKlineIntervals []string // 订阅的k线的周期列表 + subKlineInsts *types.ExchangeState[*collect.SyncMap[string, bool]] // 订阅k线中的交易产品列表 + subKlineStream grpc.BidiStreamingClient[pb.ReqStreamSubscribeKline, pb.RspStreamSubscribeKline] // 订阅k线的stream klineSignalChan chan string } -func NewKlineSeriesStore(exchangeClient pb.ExchangeServiceClient) (kss *KlineStore) { - kss = &KlineStore{ +func NewKlineSeriesStore(exchangeClient pb.ExchangeServiceClient) (kss *KlineSeriesStore) { + kss = &KlineSeriesStore{ exchangeClient: exchangeClient, - klineSignalChan: make(chan string, 128), + klineSignalChan: make(chan string, 1024), } // 周期列表 @@ -43,15 +44,15 @@ func NewKlineSeriesStore(exchangeClient pb.ExchangeServiceClient) (kss *KlineSto return collect.NewSyncMap[string, bool]() }) // 各交易所 store 初始化 - kss.store = types.NewExchangeStateInit(func() *collect.ConcurrentMap[string, *TradeInstanceKlineSeries] { - return collect.NewConcurrentMap[string, *TradeInstanceKlineSeries](64, func(s string) string { + kss.store = types.NewExchangeStateInit(func() *collect.ConcurrentMap[string, *sig.TradeInstanceKlineSeries] { + return collect.NewConcurrentMap[string, *sig.TradeInstanceKlineSeries](64, func(s string) string { return s }) }) return } -func (s *KlineStore) Init() (err error) { +func (s *KlineSeriesStore) Init() (err error) { // 连接 exchange kline stream go s.connectSubscribeKline(false) @@ -85,7 +86,7 @@ func (s *KlineStore) Init() (err error) { } // connectSubscribeKline 连接exchange订阅实时k线 -func (s *KlineStore) connectSubscribeKline(reconnect bool) { +func (s *KlineSeriesStore) connectSubscribeKline(reconnect bool) { defer func() { if s.subKlineStream != nil { s.subKlineStream.CloseSend() @@ -145,7 +146,7 @@ func (s *KlineStore) connectSubscribeKline(reconnect bool) { } // subscribeKline 发送订阅消息 -func (s *KlineStore) sendSubscribeKline(save bool, exchange pb.ExchangeType, instIds ...string) { +func (s *KlineSeriesStore) sendSubscribeKline(save bool, exchange pb.ExchangeType, instIds ...string) { if len(instIds) == 0 { return } @@ -182,13 +183,13 @@ func (s *KlineStore) sendSubscribeKline(save bool, exchange pb.ExchangeType, ins go retry.DoWithFixDelay(math.MaxInt32, time.Second, doSend) } -func (s *KlineStore) inititalKlineSeries(exchange pb.ExchangeType, instId string) { +func (s *KlineSeriesStore) inititalKlineSeries(exchange pb.ExchangeType, instId string) { if !s.store.IsSupport(exchange) { zlog.Errorf("unsupport exchange %s", exchange) return } - storeInst := s.store.Get(exchange).ComputeIfAbsent(instId, func(k string) *TradeInstanceKlineSeries { - return NewTradeInstanceKlineSeries(exchange, k) + storeInst := s.store.Get(exchange).ComputeIfAbsent(instId, func(k string) *sig.TradeInstanceKlineSeries { + return sig.NewTradeInstanceKlineSeries(exchange, k) }) // 交易产品已初始化过 if !storeInst.Status.CompareAndSwap(int32(data.StatusNone), int32(data.StatusProcessing)) { @@ -199,7 +200,7 @@ func (s *KlineStore) inititalKlineSeries(exchange pb.ExchangeType, instId string // 初始化最新的 klineSeries for _, interval := range s.subKlineIntervals { retry.DoWithFixDelay(math.MaxInt32, 2*time.Second, func(retryTimes uint32) (_ struct{}, err error) { - _, err = s.fetchHistoryKlineToSeries(exchange, instId, interval, 0, 0, MaxSeriesKlines) + _, err = s.fetchHistoryKlineToSeries(exchange, instId, interval, 0, 0, sig.MaxSeriesKlines) return }) } @@ -210,7 +211,7 @@ func (s *KlineStore) inititalKlineSeries(exchange pb.ExchangeType, instId string } // fetchHistoryKlineToSeries 拉去历史k线数据更新series -func (s *KlineStore) fetchHistoryKlineToSeries(exchange pb.ExchangeType, instId, interval string, before, after int64, count uint32) (total int, err error) { +func (s *KlineSeriesStore) fetchHistoryKlineToSeries(exchange pb.ExchangeType, instId, interval string, before, after int64, count uint32) (total int, err error) { // 拉取最新的1000条k线 req := &pb.ReqHistoryKlineStream{ Series: &pb.SeriesRange{ @@ -253,13 +254,13 @@ func (s *KlineStore) fetchHistoryKlineToSeries(exchange pb.ExchangeType, instId, return } -func (s *KlineStore) ConsumerKlineSignel() <-chan string { +func (s *KlineSeriesStore) ConsumerKlineSignel() <-chan string { return s.klineSignalChan } // Update // kline klineStore -> klineSeries -> strategy -> indicator -> klineSeries.Series -func (s *KlineStore) Update(exchange pb.ExchangeType, instId string, kline *types.Kline) { +func (s *KlineSeriesStore) Update(exchange pb.ExchangeType, instId string, kline *types.Kline) { if _, ok := types.SupportedIntervals[kline.Interval]; !ok { zlog.Warningf("unsupport interval: %s", kline.Interval) return @@ -269,8 +270,8 @@ func (s *KlineStore) Update(exchange pb.ExchangeType, instId string, kline *type return } - instSeries := s.store.Get(exchange).ComputeIfAbsent(instId, func(k string) *TradeInstanceKlineSeries { - return NewTradeInstanceKlineSeries(exchange, k) + instSeries := s.store.Get(exchange).ComputeIfAbsent(instId, func(k string) *sig.TradeInstanceKlineSeries { + return sig.NewTradeInstanceKlineSeries(exchange, k) }) before, serial := instSeries.IntervalKlines.Get(kline.Interval).Update(kline) @@ -303,27 +304,55 @@ func (s *KlineStore) Update(exchange pb.ExchangeType, instId string, kline *type return } } - // 判断同一时刻k线 + + // 判断同一时刻其它k线 var intervals []types.Interval - endTs := kline.Interval.MustAddMul(kline.Ts, 1) - instSeries.IntervalKlines.Range(func(interval types.Interval, v *KlineSeries) { - if endTs == interval.MustAddMul(v.LastTs(), 1) { + currentTs := kline.Interval.MustAddMul(kline.Ts, 1) + completed := true + instSeries.IntervalKlines.Range(func(interval types.Interval, v *sig.KlineSeries) { + lastTs := v.LastTs() + if currentTs == interval.MustAddMul(lastTs, 1) { intervals = append(intervals, interval) + } else if completed { + // 同一时刻其它k线是否未接收完成 + for i := range int64(100) { + nextTs := interval.MustAddMul(lastTs, i+2) + if currentTs == nextTs { + completed = false + } else if nextTs > currentTs { + break + } + } } }) - // zlog.Debugf("confirm kline intervals: instId=%s(%s), interval=%s, ts=%d, %v", instId, exchange, kline.Interval, kline.Ts, intervals) + zlog.Debugf("confirm kline intervals: instId=%s(%s), interval=%s, ts=%d, competed=%v, intervals=%v", instId, exchange, kline.Interval, kline.Ts, completed, intervals) // publish kline update signal - pubKey := strategy.DriverIntervalKey(instId, intervals, exchange) - select { - case s.klineSignalChan <- pubKey: - default: - zlog.Warningf("publish kline update signal fail: instId=%s(%s), interval=%s, ts=%d, %v", instId, exchange, kline.Interval, kline.Ts, intervals) + pubKeys := []string{ + strategy.DriverIntervalKey(instId, exchange, false, kline.Interval), + } + if completed { + for _, interval := range intervals { + k := strategy.DriverIntervalKey(instId, exchange, true, interval) + pubKeys = append(pubKeys, k) + } + } + // if len(intervals) > 1 { + // k := strategy.DriverIntervalKey(instId, exchange, false, intervals...) + // pubKeys = append(pubKeys, k) + // } + + for _, pubKey := range pubKeys { + select { + case s.klineSignalChan <- pubKey: + default: + zlog.Warningf("publish kline update signal fail: instId=%s(%s), interval=%s, ts=%d, %v", instId, exchange, kline.Interval, kline.Ts, intervals) + } } } // GetKlineSeires 获取k线序列 -func (s *KlineStore) GetKlineSeires(exchange pb.ExchangeType, instId string, interval types.Interval) (klineSeries *KlineSeries, err error) { +func (s *KlineSeriesStore) GetKlineSeires(exchange pb.ExchangeType, instId string, interval types.Interval) (klineSeries *sig.KlineSeries, err error) { if _, ok := types.SupportedIntervals[interval]; !ok { err = fmt.Errorf("unsupport interval: %s", interval) return diff --git a/internal/trading/libs/account_position.go b/internal/trading/libs/account_position.go deleted file mode 100644 index a8bfd5a..0000000 --- a/internal/trading/libs/account_position.go +++ /dev/null @@ -1,3 +0,0 @@ -package libs - -// 账户持仓 diff --git a/internal/trading/sig/account_position.go b/internal/trading/sig/account_position.go new file mode 100644 index 0000000..c7771e5 --- /dev/null +++ b/internal/trading/sig/account_position.go @@ -0,0 +1,3 @@ +package sig + +// 账户持仓管理 -> riskManager 风险管理 diff --git a/internal/trading/indicator_context.go b/internal/trading/sig/indicator_context.go similarity index 99% rename from internal/trading/indicator_context.go rename to internal/trading/sig/indicator_context.go index 1ccd98b..b7d7121 100644 --- a/internal/trading/indicator_context.go +++ b/internal/trading/sig/indicator_context.go @@ -1,4 +1,4 @@ -package trading +package sig import ( "context" diff --git a/internal/trading/indicator_series.go b/internal/trading/sig/indicator_series.go similarity index 98% rename from internal/trading/indicator_series.go rename to internal/trading/sig/indicator_series.go index 0091ffe..6ba7e3b 100644 --- a/internal/trading/indicator_series.go +++ b/internal/trading/sig/indicator_series.go @@ -1,4 +1,4 @@ -package trading +package sig import ( "sig-pub/pkg/indicator" diff --git a/internal/trading/indicator_service.go b/internal/trading/sig/indicator_service.go similarity index 98% rename from internal/trading/indicator_service.go rename to internal/trading/sig/indicator_service.go index ff962c6..3a6a9b9 100644 --- a/internal/trading/indicator_service.go +++ b/internal/trading/sig/indicator_service.go @@ -1,4 +1,4 @@ -package trading +package sig import ( "sig-pub/pkg/indicator" diff --git a/internal/trading/kline_series.go b/internal/trading/sig/kline_series.go similarity index 99% rename from internal/trading/kline_series.go rename to internal/trading/sig/kline_series.go index 05406b7..cb06a00 100644 --- a/internal/trading/kline_series.go +++ b/internal/trading/sig/kline_series.go @@ -1,4 +1,4 @@ -package trading +package sig import ( "fmt" diff --git a/internal/trading/strategy_context.go b/internal/trading/sig/strategy_context.go similarity index 56% rename from internal/trading/strategy_context.go rename to internal/trading/sig/strategy_context.go index 85c6da1..3825df1 100644 --- a/internal/trading/strategy_context.go +++ b/internal/trading/sig/strategy_context.go @@ -1,4 +1,4 @@ -package trading +package sig import ( "fmt" @@ -39,37 +39,6 @@ func (c *StrategyContext) Series(offset, count int16) (klines series.Klines) { return c.indicatorContext.Series(offset, count) } -// // Buy 发出多信号 -// func (c *StrategyContext) Buy() { -// zlog.Infof("signal buy: %d", c.Get(0).Ts) - -// c.signal = append(c.signal, pb.Side_BUY) -// c.signalTimes = append(c.signalTimes, c.Get(0).Ts) -// win := false -// signalPrice := c.Get(0).Close -// if c.indicatorContext.GetOffset() > 0 { -// c.indicatorContext.AddOffset(-1) -// win = c.Get(0).Close.Cmp(signalPrice) > 0 -// c.indicatorContext.AddOffset(1) -// } -// c.wins = append(c.wins, win) -// } - -// // Sell 发出空信号 -// func (c *StrategyContext) Sell() { -// zlog.Infof("signal sell: %d", c.Get(0).Ts) -// c.signal = append(c.signal, pb.Side_SELL) -// c.signalTimes = append(c.signalTimes, c.Get(0).Ts) -// win := false -// signalPrice := c.Get(0).Close -// if c.indicatorContext.GetOffset() > 0 { -// c.indicatorContext.AddOffset(-1) -// win = c.Get(0).Close.Cmp(signalPrice) < 0 -// c.indicatorContext.AddOffset(1) -// } -// c.wins = append(c.wins, win) -// } - // 获取窗口类型指标 func (c *StrategyContext) IndicatorW(name string, window int16) (s indicator.IIndicatorSeries) { indicator, ok := c.indicatorsReg.IndicatorW(name) @@ -78,3 +47,33 @@ func (c *StrategyContext) IndicatorW(name string, window int16) (s indicator.IIn } return NewWindowIndicatorSeries(window, indicator, c.indicatorContext) } + +type IOffsetIntervalStrategyContext interface { + strategy.IIntervalStrategyContext + SetOffset(offset int16) +} + +// IntervalStrategyContext 周期策略上下文 +type IntervalStrategyContext struct { + IOffsetIntervalStrategyContext +} + +// Get [0]当前k线 +func (c *IntervalStrategyContext) Get(interval types.Interval, offset int16) (kline types.Kline) { + return +} + +// Series [offset...end] +func (c *IntervalStrategyContext) Series(interval types.Interval, offset, count int16) (klines series.Klines) { + return +} + +// 获取窗口类型指标 +func (c *IntervalStrategyContext) IndicatorW(interval types.Interval, name string, window int16) (series indicator.IIndicatorSeries) { + return +} + +// 获取其它策略 +// func (c *IntervalStrategyContext) SigStrategy(interval types.Interval, name string, sigParam strategy.SigStrategyParam) (sigStrategy strategy.ISigStrategy) { +// return +// } diff --git a/internal/trading/trading_plan_runner.go b/internal/trading/sig/trading_plan_runner.go similarity index 98% rename from internal/trading/trading_plan_runner.go rename to internal/trading/sig/trading_plan_runner.go index e501d9c..86e1c27 100644 --- a/internal/trading/trading_plan_runner.go +++ b/internal/trading/sig/trading_plan_runner.go @@ -1,4 +1,4 @@ -package trading +package sig import ( "sig-pub/pkg/data/entity" diff --git a/internal/trading/trading_service.go b/internal/trading/trading_service.go index 3672e05..d745a4f 100644 --- a/internal/trading/trading_service.go +++ b/internal/trading/trading_service.go @@ -13,6 +13,8 @@ import ( "sig-pub/pkg/utils/collect" "sig-pub/pkg/zlog" + "sig-pub/internal/trading/sig" + "github.com/bytedance/sonic" ) @@ -20,11 +22,11 @@ type TradingService struct { marketClientAside *client.TradeInstanceAside exchangeClient pb.ExchangeServiceClient - klineStore *KlineStore - indicatorReg *indicator.IndicatorRegistry // 注册窗口指标 - strategyReg *strategy.SigStrategyRegistry // 注册信号策略 - signalPublisher *publish.Publisher[int64, strategy.StrategyType] // planId -> strategyType - tradingPlans *collect.SyncMap[int64, *TradingPlan] // 运行中交易计划 + klineSeriesStore *KlineSeriesStore + indicatorReg *indicator.IndicatorRegistry // 注册窗口指标 + strategyReg *strategy.SigStrategyRegistry // 注册信号策略 + signalPublisher *publish.Publisher[int64, strategy.StrategyType] // planId -> strategyType + tradingPlans *collect.SyncMap[int64, *sig.TradingPlan] // 运行中交易计划 } func NewTradingService( @@ -34,11 +36,11 @@ func NewTradingService( return &TradingService{ marketClientAside: marketClientAside, exchangeClient: exchangeClient, - klineStore: NewKlineSeriesStore(exchangeClient), + klineSeriesStore: NewKlineSeriesStore(exchangeClient), indicatorReg: indicator.NewIndicatorRegistry(), strategyReg: strategy.NewSigStrategyRegistry(), signalPublisher: publish.NewPublisher[int64, strategy.StrategyType](16), - tradingPlans: collect.NewSyncMap[int64, *TradingPlan](), + tradingPlans: collect.NewSyncMap[int64, *sig.TradingPlan](), } } @@ -51,7 +53,7 @@ func (svc *TradingService) Init() (err error) { return } - if err = svc.klineStore.Init(); err != nil { + if err = svc.klineSeriesStore.Init(); err != nil { return } @@ -63,7 +65,7 @@ func (svc *TradingService) Init() (err error) { // consumerKlineSignal 订阅k线更新 func (svc *TradingService) consumerKlineSignal() { - c := svc.klineStore.ConsumerKlineSignel() + c := svc.klineSeriesStore.ConsumerKlineSignel() for { signalKey := <-c zlog.Debugf("signal: %s", signalKey) @@ -96,7 +98,7 @@ func (svc *TradingService) runTradingPlan(plan *entity.TradePlan) (err error) { return } - tradingPlan := NewTradingPlan(*plan, svc.indicatorReg) + tradingPlan := sig.NewTradingPlan(*plan, svc.indicatorReg) if _, load := svc.tradingPlans.LoadOrStore(planId, tradingPlan); load { err = fmt.Errorf("plan already running: planId=%d", planId) return @@ -108,7 +110,7 @@ func (svc *TradingService) runTradingPlan(plan *entity.TradePlan) (err error) { svc.tradingPlans.Delete(planId) } else { // 订阅交易信号策略k线周期 - sigSubKey := strategy.DriverIntervalKey(instId, []types.Interval{sigInterval}, exchange) + sigSubKey := strategy.DriverIntervalKey(instId, exchange, false, sigInterval) svc.signalPublisher.Subscribe(sigSubKey, planId, strategy.StrategyTypeSig) tradingPlan.Status.Store(int32(data.StatusOk)) @@ -130,7 +132,7 @@ func (svc *TradingService) runTradingPlan(plan *entity.TradePlan) (err error) { err = fmt.Errorf("strategy %s not exists", plan.SigStrategy) return } - sigKlineSeries, err := svc.klineStore.GetKlineSeires(exchange, instId, sigInterval) + sigKlineSeries, err := svc.klineSeriesStore.GetKlineSeires(exchange, instId, sigInterval) if err != nil { return } @@ -138,7 +140,7 @@ func (svc *TradingService) runTradingPlan(plan *entity.TradePlan) (err error) { if err = tradingPlan.Init(); err != nil { return } - sigIndCtx := NewIndicatorContext(sigKlineSeries) + sigIndCtx := sig.NewIndicatorContext(sigKlineSeries) if err = tradingPlan.InitSigStrategy(sigStrategy, *sigStrategyParam, sigIndCtx); err != nil { return } @@ -175,7 +177,7 @@ func (svc *TradingService) IndicatorSeries(indicatorName string, window uint32, // 查询历史指标数据 sr.Window = window - indCtx := NewHistoryIndicatorContext(svc.exchangeClient) + indCtx := sig.NewHistoryIndicatorContext(svc.exchangeClient) totalK := 0 if totalK, err = indCtx.Init(sr); err != nil { return @@ -227,14 +229,14 @@ func (svc *TradingService) StrategySeries(req *pb.ReqStrategySeries, rsp *pb.Rsp // } // recover todo out of range count, totalK := 0, 0 - indicatorContext := NewHistoryIndicatorContext(svc.exchangeClient) + indicatorContext := sig.NewHistoryIndicatorContext(svc.exchangeClient) req.Series.Window += MaxIndicatorWindow if totalK, err = indicatorContext.Init(req.Series); err != nil { return } count = totalK - MaxIndicatorWindow - strategyContext := NewStrategyContext(indicatorContext, svc.indicatorReg) + strategyContext := sig.NewStrategyContext(indicatorContext, svc.indicatorReg) for i := count - 1; i >= 0; i-- { strategyContext.SetOffset(int16(i)) side := sigStrategy.Update(strategyContext) diff --git a/pkg/grpc/client/direct_client_factory.go b/pkg/grpc/client/direct_client_factory.go index 6f6b2ce..6ba88b7 100755 --- a/pkg/grpc/client/direct_client_factory.go +++ b/pkg/grpc/client/direct_client_factory.go @@ -2,8 +2,9 @@ package client import ( "context" - "google.golang.org/grpc" "sync" + + "google.golang.org/grpc" ) type GrpcDirectClientFactory struct { @@ -22,8 +23,8 @@ func (f *GrpcDirectClientFactory) NewConn(ctx context.Context, addr string, opts dialOpts := make([]grpc.DialOption, 0, len(f.defaultOpts)+len(opts)) dialOpts = append(dialOpts, f.defaultOpts...) dialOpts = append(dialOpts, opts...) - - return grpc.DialContext(ctx, addr, dialOpts...) + return grpc.NewClient(addr, dialOpts...) + // return grpc.DialContext(ctx, addr, dialOpts...) } func (f *GrpcDirectClientFactory) GetConn(ctx context.Context, addr string, opts ...grpc.DialOption) (conn *grpc.ClientConn, err error) { diff --git a/pkg/strategy/gold_x.go b/pkg/strategy/gold_x.go index d945680..59ed395 100644 --- a/pkg/strategy/gold_x.go +++ b/pkg/strategy/gold_x.go @@ -3,11 +3,13 @@ package strategy import ( "fmt" "sig-pub/api/pb" + "sig-pub/pkg/types" ) // GoldX 金叉策略 type GoldX struct { ISigStrategy + IIntervalSigStrategy short, long int16 } @@ -56,3 +58,9 @@ func (s *GoldX) Update(ctx ISigStrategyContext) (side pb.Side) { } return } + +func (s *GoldX) UpdateByIntervals(ctx IIntervalStrategyContext) (side pb.Side) { + series5m := ctx.Series(types.Interval5m, 0, 2) + series5m.Close().Diff() + return +} diff --git a/pkg/strategy/sig_strategy.go b/pkg/strategy/sig_strategy.go index a6e2dcf..36c2a38 100644 --- a/pkg/strategy/sig_strategy.go +++ b/pkg/strategy/sig_strategy.go @@ -1,13 +1,10 @@ package strategy import ( - "fmt" "sig-pub/api/pb" "sig-pub/pkg/indicator" "sig-pub/pkg/types" "sig-pub/pkg/types/series" - "sig-pub/pkg/utils/collect" - "strings" ) // ISigStrategy 交易信号策略接口(单周期单交易所) @@ -34,14 +31,3 @@ type ISigStrategyContext interface { // 获取窗口类型指标 IndicatorW(name string, window int16) indicator.IIndicatorSeries } - -// DriverIntervalKey 生成周期驱动事件key -// interval/BTC_USDT/OKX,BINANCE/1m,3m,5m -func DriverIntervalKey(instId string, intervals []types.Interval, exchanges ...pb.ExchangeType) string { - types.IntervalsSort(intervals) - types.ExchangesSort(exchanges) - strIntervals := collect.Mapping(intervals, func(_ int, interval types.Interval) string { return string(interval) }) - strExchanges := collect.Mapping(exchanges, func(_ int, exchange pb.ExchangeType) string { return exchange.String() }) - pubKey := fmt.Sprintf("/interval/%s/%s/%s", instId, strings.Join(strExchanges, ","), strings.Join(strIntervals, ",")) - return pubKey -} diff --git a/pkg/strategy/sig_strategy_exchanges.go b/pkg/strategy/sig_strategy_exchanges.go index ab8a6e1..64add10 100644 --- a/pkg/strategy/sig_strategy_exchanges.go +++ b/pkg/strategy/sig_strategy_exchanges.go @@ -1,5 +1,11 @@ package strategy +import ( + "sig-pub/pkg/indicator" + "sig-pub/pkg/types" + "sig-pub/pkg/types/series" +) + // IExchangesSigStrategy 多交易所策略 type IExchangesSigStrategy interface { ISigStrategy @@ -9,3 +15,17 @@ type IExchangesSigStrategy interface { type IExchangesIntervalsStrategy interface { ISigStrategy } + +// IExchangesSigStrategyContext 策略外部访问能力 +// klineSeries, Indicator +type IExchangesSigStrategyContext interface { + Buy() // 发出多信号 + Sell() // 发出空信号 + + // Get [0]当前k线 + Get(interval types.Interval, offset int16) types.Kline + // Series [offset...end] + Series(interval types.Interval, offset, count int16) (klines series.Klines) + // 获取窗口类型指标 + IndicatorW(interval types.Interval, name string, window int16) indicator.IIndicatorSeries +} diff --git a/pkg/strategy/sig_strategy_intervals.go b/pkg/strategy/sig_strategy_intervals.go index ffae5f9..f97189b 100644 --- a/pkg/strategy/sig_strategy_intervals.go +++ b/pkg/strategy/sig_strategy_intervals.go @@ -1,26 +1,19 @@ package strategy import ( + "sig-pub/api/pb" "sig-pub/pkg/indicator" "sig-pub/pkg/types" "sig-pub/pkg/types/series" ) -// 多周期k线策略 -type IIntervalsSigStrategy interface { - New() IIntervalsSigStrategy - Meta() StrategyMeta - Update(ctx IIntervalsSigStrategyContext) - DriverIntervals() []types.Interval // 驱动k线周期, 当驱动周期k线更新时则判断调用Update方法 - SubscribeIntervals() []types.Interval // 订阅k线周期, 当同一时间的订阅周期都更新时调用Update方法 +// 多周期k线策略接口 +type IIntervalSigStrategy interface { + UpdateByInterval(ctx IIntervalStrategyContext) (side pb.Side) } -// ISigStrategyContext 策略外部访问能力 -// klineSeries, Indicator -type IIntervalsSigStrategyContext interface { - Buy() // 发出多信号 - Sell() // 发出空信号 - +// IIntervalStrategyContext 多周期策略上下文 +type IIntervalStrategyContext interface { // Get [0]当前k线 Get(interval types.Interval, offset int16) types.Kline // Series [offset...end] diff --git a/pkg/strategy/sig_strategy_registry.go b/pkg/strategy/sig_strategy_registry.go index 52ea496..da00444 100644 --- a/pkg/strategy/sig_strategy_registry.go +++ b/pkg/strategy/sig_strategy_registry.go @@ -7,14 +7,14 @@ import ( // 指标注册器 type SigStrategyRegistry struct { - sigStrategies *collect.SyncMap[string, ISigStrategy] // 注册信号策略 - intervalSigStrategies *collect.SyncMap[string, IIntervalsSigStrategy] // 注册窗口指标 + sigStrategies *collect.SyncMap[string, ISigStrategy] // 注册信号策略 + intervalSigStrategies *collect.SyncMap[string, IIntervalSigStrategy] // 注册窗口指标 } func NewSigStrategyRegistry() *SigStrategyRegistry { return &SigStrategyRegistry{ sigStrategies: collect.NewSyncMap[string, ISigStrategy](), - intervalSigStrategies: collect.NewSyncMap[string, IIntervalsSigStrategy](), + intervalSigStrategies: collect.NewSyncMap[string, IIntervalSigStrategy](), } } diff --git a/pkg/strategy/strategy.go b/pkg/strategy/strategy.go index b6b5fb3..b90cc4b 100644 --- a/pkg/strategy/strategy.go +++ b/pkg/strategy/strategy.go @@ -1,5 +1,14 @@ package strategy +import ( + "fmt" + "sig-pub/api/pb" + "sig-pub/pkg/types" + "sig-pub/pkg/utils/collect" + "sig-pub/pkg/utils/lang" + "strings" +) + type StrategyType int32 const ( @@ -8,3 +17,24 @@ const ( StrategyTypeTrade // 交易下单策略 StrategyTypeClose // 交易平仓策略 ) + +// DriverIntervalKey 生成周期驱动事件key +// interval/BTC_USDT/OKX/1m,3m,5m/1 +// completed 同一时刻的所有其他周期都完成 +func DriverIntervalKey(instId string, exchange pb.ExchangeType, completed bool, intervals ...types.Interval) string { + types.IntervalsSort(intervals) + strIntervals := collect.Mapping(intervals, func(_ int, interval types.Interval) string { return string(interval) }) + pubKey := fmt.Sprintf("/interval/%s/%s/%s/%d", instId, exchange.String(), strings.Join(strIntervals, ","), lang.Ternary(completed, 1, 0)) + return pubKey +} + +// DriverIntervalKey 生成周期驱动事件key +// interval/BTC_USDT/OKX,BINANCE/1m,3m,5m +func MulExchangeDriverIntervalKey(instId string, exchanges []pb.ExchangeType, intervals ...types.Interval) string { + types.IntervalsSort(intervals) + types.ExchangesSort(exchanges) + strIntervals := collect.Mapping(intervals, func(_ int, interval types.Interval) string { return string(interval) }) + strExchanges := collect.Mapping(exchanges, func(_ int, exchange pb.ExchangeType) string { return exchange.String() }) + pubKey := fmt.Sprintf("/interval/%s/%s/%s", instId, strings.Join(strExchanges, ","), strings.Join(strIntervals, ",")) + return pubKey +} diff --git a/pkg/utils/lang/condition.go b/pkg/utils/lang/condition.go new file mode 100644 index 0000000..5d4e771 --- /dev/null +++ b/pkg/utils/lang/condition.go @@ -0,0 +1,19 @@ +package lang + +// Ternary is a 1 line if/else statement. +func Ternary[T any](condition bool, ifOutput T, elseOutput T) T { + if condition { + return ifOutput + } + + return elseOutput +} + +// TernaryF is a 1 line if/else statement whose options are functions +func TernaryF[T any](condition bool, ifFunc func() T, elseFunc func() T) T { + if condition { + return ifFunc() + } + + return elseFunc() +}