From e9431c85593eee11deb06bf1444d65a9f9f16fe0 Mon Sep 17 00:00:00 2001 From: strange Date: Fri, 24 Oct 2025 18:36:50 +0800 Subject: [PATCH] trading plan --- internal/trading/trading_plan_runner.go | 49 +++++++------ internal/trading/trading_service.go | 94 ++++++++++++++++++++----- pkg/data/entity/trade_plan.go | 28 ++++---- pkg/strategy/strategy.go | 10 +++ pkg/utils/collect/sync_map.go | 4 ++ 5 files changed, 130 insertions(+), 55 deletions(-) create mode 100644 pkg/strategy/strategy.go diff --git a/internal/trading/trading_plan_runner.go b/internal/trading/trading_plan_runner.go index e3a49fb..8086946 100644 --- a/internal/trading/trading_plan_runner.go +++ b/internal/trading/trading_plan_runner.go @@ -1,40 +1,38 @@ package trading import ( - "fmt" "sig-pub/pkg/data/entity" "sig-pub/pkg/indicator" - "sig-pub/pkg/publish" "sig-pub/pkg/strategy" + "sync/atomic" ) -type TradingPlanRunner struct { +type TradingPlan struct { + Status atomic.Int32 plan entity.TradePlan - indicatorReg indicator.IndicatorRegistry - strategyReg strategy.SigStrategyRegistry - publisher publish.Publisher[int32, any] + indicatorReg *indicator.IndicatorRegistry + + sigStrategy strategy.ISigStrategy + sigStrategyContext *StrategyContext + // publisher publish.Publisher[int32, any] + // signalKey map[string]int32 } -func NewTradingPlan(plan entity.TradePlan, - indicatorReg indicator.IndicatorRegistry, - strategyReg strategy.SigStrategyRegistry, -) *TradingPlanRunner { - return &TradingPlanRunner{ +func NewTradingPlan(plan entity.TradePlan, indicatorReg *indicator.IndicatorRegistry) *TradingPlan { + return &TradingPlan{ plan: plan, indicatorReg: indicatorReg, - strategyReg: strategyReg, } } // Init 初始化交易计划 // subKlineKeys 订阅k线更新, 更新时调用Update方法 // 交易信号/下单/平仓/风控 -func (r *TradingPlanRunner) Init() (subSignalKeyKeys []string, err error) { +func (r *TradingPlan) Init() (err error) { // 初始化执行策略 - _, err = r.initSigStrategy(r.plan.SigStrategy, r.plan.SigStrategyParams) - if err != nil { - return - } + // if err = r.initSigStrategy(r.sigStrategy, r.plan.SigStrategyParams); err != nil { + // return + // } // sigStrategy 1m // tradeStrategy 1s @@ -44,18 +42,19 @@ func (r *TradingPlanRunner) Init() (subSignalKeyKeys []string, err error) { // initSigStrategy 初始化多空信号策略 // buy/sell -> 过滤/风控 -> tradeStrategy -> closeStrategy -func (r *TradingPlanRunner) initSigStrategy(name, params string) (subSignalKeyKeys []string, err error) { - strategy, ok := r.strategyReg.NewSigStrategy(name) - if !ok { - err = fmt.Errorf("strategy not exists: %s", name) +func (r *TradingPlan) InitSigStrategy(sigStrategy strategy.ISigStrategy, params strategy.SigStrategyParam, klineSeries *KlineSeries) (err error) { + if err = sigStrategy.Init(params); err != nil { return } - _ = strategy - // strategy.Update() + r.sigStrategy = sigStrategy + r.sigStrategyContext = NewStrategyContext(klineSeries, r.indicatorReg) return } // Update 订阅k线更新 -func (r *TradingPlanRunner) Update(signalKey string) { - +func (r *TradingPlan) Update(signalType strategy.StrategyType) { + switch signalType { + case strategy.StrategyTypeSig: + r.sigStrategy.Update(r.sigStrategyContext) + } } diff --git a/internal/trading/trading_service.go b/internal/trading/trading_service.go index cbfa38e..222fabc 100644 --- a/internal/trading/trading_service.go +++ b/internal/trading/trading_service.go @@ -4,6 +4,7 @@ import ( "fmt" "sig-pub/api/pb" "sig-pub/pkg/client" + "sig-pub/pkg/data" "sig-pub/pkg/data/entity" "sig-pub/pkg/indicator" "sig-pub/pkg/publish" @@ -11,17 +12,19 @@ import ( "sig-pub/pkg/types" "sig-pub/pkg/utils/collect" "sig-pub/pkg/zlog" + + "github.com/bytedance/sonic" ) type TradingService struct { marketClientAside *client.TradeInstanceAside exchangeClient pb.ExchangeServiceClient - klineStore *KlineStore - indicatorReg *indicator.IndicatorRegistry // 注册窗口指标 - strategyReg *strategy.SigStrategyRegistry // 注册信号策略 - publisher *publish.Publisher[int64, *TradingPlanRunner] - tradingPlan chan *TradingPlanRunner + klineStore *KlineStore + indicatorReg *indicator.IndicatorRegistry // 注册窗口指标 + strategyReg *strategy.SigStrategyRegistry // 注册信号策略 + signalPublisher *publish.Publisher[int64, strategy.StrategyType] // planId -> strategyType + tradingPlans *collect.SyncMap[int64, *TradingPlan] // 运行中交易计划 } func NewTradingService( @@ -34,7 +37,8 @@ func NewTradingService( klineStore: NewKlineSeriesStore(exchangeClient), indicatorReg: indicator.NewIndicatorRegistry(), strategyReg: strategy.NewSigStrategyRegistry(), - publisher: publish.NewPublisher[int64, *TradingPlanRunner](8), + signalPublisher: publish.NewPublisher[int64, strategy.StrategyType](16), + tradingPlans: collect.NewSyncMap[int64, *TradingPlan](), } } @@ -52,6 +56,8 @@ func (svc *TradingService) Init() (err error) { } go svc.consumerKlineSignal() + + // todo loading trading plan return } @@ -61,24 +67,80 @@ func (svc *TradingService) consumerKlineSignal() { for { signalKey := <-c zlog.Debugf("signal: %s", signalKey) - _, plans := svc.publisher.Publisher(signalKey) - for _, plan := range plans { - plan.Update(signalKey) + planIds, strategyTypes := svc.signalPublisher.Publisher(signalKey) + for i, strategyType := range strategyTypes { + planId := planIds[i] + plan, ok := svc.tradingPlans.Load(planId) + if !ok { + zlog.Warningf("plan not running: id=%d", planId) + continue + } + if plan.Status.Load() == int32(data.StatusOk) { + plan.Update(strategyType) + } } } } -// RunStrategy 运行策略 -// todo 止盈止损... -func (svc *TradingService) RunQuantPlan(plan *entity.TradePlan) (err error) { - strategy, ok := svc.strategyReg.NewSigStrategy(plan.SigStrategy) +// runTradingPlan 运行交易计划 +// todo 止盈止损策略, 下单策略... +func (svc *TradingService) runTradingPlan(plan *entity.TradePlan) (err error) { + var planId = plan.Id + var instId = plan.InstId + var exchange pb.ExchangeType + var sigInterval types.Interval + + exchange = pb.ExchangeType(plan.Exchange) + if !types.IsSupportExchange(exchange) { + err = fmt.Errorf("unsupport exchange %d", plan.Exchange) + return + } + + tradingPlan := NewTradingPlan(*plan, svc.indicatorReg) + if _, load := svc.tradingPlans.LoadOrStore(planId, tradingPlan); load { + err = fmt.Errorf("plan already running: planId=%d", planId) + return + } + tradingPlan.Status.Store(int32(data.StatusProcessing)) + + defer func() { + if err != nil { + svc.tradingPlans.Delete(planId) + } else { + // 订阅交易信号策略k线周期 + sigSubKey := strategy.DriverIntervalKey(instId, []types.Interval{sigInterval}, exchange) + svc.signalPublisher.Subscribe(sigSubKey, planId, strategy.StrategyTypeSig) + + tradingPlan.Status.Store(int32(data.StatusOk)) + } + }() + + // sigStrategy + sigStrategyParam := new(strategy.SigStrategyParam) + if err = sonic.UnmarshalString(plan.SigStrategyParam, sigStrategyParam); err != nil { + return + } + sigInterval = types.Interval(sigStrategyParam.Interval) + if _, ok := types.SupportedIntervals[sigInterval]; !ok { + err = fmt.Errorf("unsupport interval %d", plan.Exchange) + return + } + sigStrategy, ok := svc.strategyReg.NewSigStrategy(plan.SigStrategy) if !ok { err = fmt.Errorf("strategy %s not exists", plan.SigStrategy) return } - runner := strategy.New() - _ = runner - runner.Update(nil) + sigKlineSeries, err := svc.klineStore.GetKlineSeires(exchange, instId, sigInterval) + if err != nil { + return + } + + if err = tradingPlan.Init(); err != nil { + return + } + if err = tradingPlan.InitSigStrategy(sigStrategy, *sigStrategyParam, sigKlineSeries); err != nil { + return + } return } diff --git a/pkg/data/entity/trade_plan.go b/pkg/data/entity/trade_plan.go index 5ac0a19..70f3f62 100644 --- a/pkg/data/entity/trade_plan.go +++ b/pkg/data/entity/trade_plan.go @@ -2,20 +2,20 @@ package entity // TradePlan 交易计划 type TradePlan struct { - Id int64 `gorm:"column:id;primaryKey" json:"id"` // id - UserId int64 `gorm:"column:user_id" json:"userId"` // 用户id - Status int8 `gorm:"column:status" json:"status"` // 状态:0禁用,1启用 - Exchange int8 `gorm:"column:exchange" json:"exchange"` // 交易所 - InstId string `gorm:"column:inst_id" json:"instId"` // 交易产品id - Interval string `gorm:"column:interval" json:"interval"` // 交易周期 - SigStrategy string `gorm:"column:sig_strategy" json:"sigStrategy"` // 交易信号策略 - ExitStrategy string `gorm:"column:exit_strategy" json:"exitStrategy"` // 退出策略 - TradeStrategy string `gorm:"column:trade_strategy" json:"tradeStrategy"` // 下单仓位管理策略 - SigStrategyParams string `gorm:"column:sig_strategy_params" json:"sigStrategyParams"` // 交易信号策略参数 - ExitStrategyParams string `gorm:"column:exit_strategy_params" json:"exitStrategyParams"` // 退出策略名称参数 - TradeStrategyParams string `gorm:"column:trade_strategy_params" json:"tradeStrategyParams"` // 下单仓位管理策略参数 - UpdateBy string `gorm:"column:update_by" json:"updateBy"` // 更新人 - UpdateTime int64 `gorm:"column:update_time" json:"updateTime"` // 更新时间戳毫秒 + Id int64 `gorm:"column:id;primaryKey" json:"id"` // id + UserId int64 `gorm:"column:user_id" json:"userId"` // 用户id + Status int8 `gorm:"column:status" json:"status"` // 状态:0禁用,1启用 + Exchange int8 `gorm:"column:exchange" json:"exchange"` // 交易所 + InstId string `gorm:"column:inst_id" json:"instId"` // 交易产品id + Interval string `gorm:"column:interval" json:"interval"` // 交易周期 + SigStrategy string `gorm:"column:sig_strategy" json:"sigStrategy"` // 交易信号策略 + ExitStrategy string `gorm:"column:exit_strategy" json:"exitStrategy"` // 退出策略 + TradeStrategy string `gorm:"column:trade_strategy" json:"tradeStrategy"` // 下单仓位管理策略 + SigStrategyParam string `gorm:"column:sig_strategy_param" json:"sigStrategyParam"` // 交易信号策略参数 + ExitStrategyParam string `gorm:"column:exit_strategy_param" json:"exitStrategyParam"` // 退出策略名称参数 + TradeStrategyParam string `gorm:"column:trade_strategy_param" json:"tradeStrategyParam"` // 下单仓位管理策略参数 + UpdateBy string `gorm:"column:update_by" json:"updateBy"` // 更新人 + UpdateTime int64 `gorm:"column:update_time" json:"updateTime"` // 更新时间戳毫秒 } func (TradePlan) TableName() string { diff --git a/pkg/strategy/strategy.go b/pkg/strategy/strategy.go new file mode 100644 index 0000000..b6b5fb3 --- /dev/null +++ b/pkg/strategy/strategy.go @@ -0,0 +1,10 @@ +package strategy + +type StrategyType int32 + +const ( + _ StrategyType = iota + StrategyTypeSig // 交易信号策略 + StrategyTypeTrade // 交易下单策略 + StrategyTypeClose // 交易平仓策略 +) diff --git a/pkg/utils/collect/sync_map.go b/pkg/utils/collect/sync_map.go index 7cf2e4d..1655093 100644 --- a/pkg/utils/collect/sync_map.go +++ b/pkg/utils/collect/sync_map.go @@ -27,6 +27,10 @@ func (m *SyncMap[K, V]) Load(k K) (v V, ok bool) { return } +func (m *SyncMap[K, V]) Delete(k K) { + m.m.Delete(k) +} + func (m *SyncMap[K, V]) Range(f func(k K, v V) bool) { m.m.Range(func(key, value any) bool { return f(key.(K), value.(V))