bbgo_origin/pkg/bbgo/environment.go

375 lines
11 KiB
Go
Raw Normal View History

2020-10-16 02:14:36 +00:00
package bbgo
import (
"context"
2020-10-20 04:11:44 +00:00
"fmt"
2020-10-16 02:14:36 +00:00
"strings"
"time"
"github.com/jmoiron/sqlx"
2020-10-20 04:11:44 +00:00
"github.com/pkg/errors"
2020-10-17 16:06:08 +00:00
log "github.com/sirupsen/logrus"
2020-10-16 02:14:36 +00:00
2020-10-30 21:21:17 +00:00
"github.com/c9s/bbgo/pkg/accounting/pnl"
2020-10-16 02:14:36 +00:00
"github.com/c9s/bbgo/pkg/service"
"github.com/c9s/bbgo/pkg/types"
2020-10-30 21:21:17 +00:00
"github.com/c9s/bbgo/pkg/util"
2020-10-16 02:14:36 +00:00
)
var LoadedExchangeStrategies = make(map[string]SingleExchangeStrategy)
var LoadedCrossExchangeStrategies = make(map[string]CrossExchangeStrategy)
2020-10-28 23:54:59 +00:00
func RegisterStrategy(key string, s interface{}) {
switch d := s.(type) {
case SingleExchangeStrategy:
LoadedExchangeStrategies[key] = d
case CrossExchangeStrategy:
LoadedCrossExchangeStrategies[key] = d
}
}
2020-10-20 05:52:25 +00:00
var emptyTime time.Time
2020-10-16 02:14:36 +00:00
// Environment presents the real exchange data layer
type Environment struct {
// Notifiability here for environment is for the streaming data notification
// note that, for back tests, we don't need notification.
Notifiability
2020-10-16 02:14:36 +00:00
TradeService *service.TradeService
TradeSync *service.SyncService
2020-10-16 02:14:36 +00:00
// startTime is the time of start point (which is used in the backtest)
startTime time.Time
2020-10-20 05:11:04 +00:00
tradeScanTime time.Time
sessions map[string]*ExchangeSession
2020-10-16 02:14:36 +00:00
}
func NewEnvironment() *Environment {
2020-10-16 02:14:36 +00:00
return &Environment{
2020-10-26 05:48:59 +00:00
// default trade scan time
2020-10-20 05:11:04 +00:00
tradeScanTime: time.Now().AddDate(0, 0, -7), // sync from 7 days ago
sessions: make(map[string]*ExchangeSession),
2020-10-16 02:14:36 +00:00
}
}
2020-10-26 05:48:59 +00:00
func (environ *Environment) SyncTrades(db *sqlx.DB) *Environment {
environ.TradeService = &service.TradeService{DB: db}
environ.TradeSync = &service.SyncService{
TradeService: environ.TradeService,
2020-10-26 05:48:59 +00:00
}
return environ
}
2020-10-16 02:14:36 +00:00
func (environ *Environment) AddExchange(name string, exchange types.Exchange) (session *ExchangeSession) {
2020-10-17 15:51:44 +00:00
session = NewExchangeSession(name, exchange)
2020-10-16 02:14:36 +00:00
environ.sessions[name] = session
return session
}
// Init prepares the data that will be used by the strategies
2020-10-16 02:14:36 +00:00
func (environ *Environment) Init(ctx context.Context) (err error) {
for n := range environ.sessions {
var session = environ.sessions[n]
2020-10-20 04:11:44 +00:00
var markets types.MarketMap
err = WithCache(fmt.Sprintf("%s-markets", session.Exchange.Name()), &markets, func() (interface{}, error) {
2020-10-20 04:11:44 +00:00
return session.Exchange.QueryMarkets(ctx)
})
2020-10-16 02:14:36 +00:00
if err != nil {
return err
}
2020-10-20 04:11:44 +00:00
if len(markets) == 0 {
return errors.Errorf("market config should not be empty")
}
2020-10-18 04:30:13 +00:00
session.markets = markets
2020-10-19 13:58:50 +00:00
2020-10-19 14:02:05 +00:00
// trade sync and market data store depends on subscribed symbols so we have to do this here.
for symbol := range session.loadedSymbols {
2020-10-16 02:14:36 +00:00
var trades []types.Trade
2020-10-26 05:48:59 +00:00
if environ.TradeSync != nil {
log.Infof("syncing trades from %s for symbol %s...", session.Exchange.Name(), symbol)
if err := environ.TradeSync.SyncTrades(ctx, session.Exchange, symbol, environ.tradeScanTime); err != nil {
2020-10-26 05:48:59 +00:00
return err
}
2020-10-16 02:14:36 +00:00
2020-10-26 05:48:59 +00:00
tradingFeeCurrency := session.Exchange.PlatformFeeCurrency()
if strings.HasPrefix(symbol, tradingFeeCurrency) {
trades, err = environ.TradeService.QueryForTradingFeeCurrency(session.Exchange.Name(), symbol, tradingFeeCurrency)
2020-10-26 05:48:59 +00:00
} else {
trades, err = environ.TradeService.Query(session.Exchange.Name(), symbol)
2020-10-26 05:48:59 +00:00
}
if err != nil {
return err
}
log.Infof("symbol %s: %d trades loaded", symbol, len(trades))
2020-10-16 02:14:36 +00:00
}
session.Trades[symbol] = trades
session.lastPrices[symbol] = 0.0
2020-10-17 16:06:08 +00:00
marketDataStore := NewMarketDataStore(symbol)
2020-10-19 13:58:50 +00:00
marketDataStore.BindStream(session.Stream)
2020-10-28 01:43:19 +00:00
session.marketDataStores[symbol] = marketDataStore
standardIndicatorSet := NewStandardIndicatorSet(symbol, marketDataStore)
session.standardIndicatorSets[symbol] = standardIndicatorSet
2020-10-28 01:43:19 +00:00
}
2020-10-28 01:13:57 +00:00
log.Infof("querying balances...")
balances, err := session.Exchange.QueryAccountBalances(ctx)
if err != nil {
return err
}
session.Account.UpdateBalances(balances)
session.Account.BindStream(session.Stream)
// update last prices
session.Stream.OnKLineClosed(func(kline types.KLine) {
log.Infof("kline closed: %+v", kline)
session.lastPrices[kline.Symbol] = kline.Close
session.marketDataStores[kline.Symbol].AddKLine(kline)
})
session.Stream.OnTradeUpdate(func(trade types.Trade) {
session.Trades[trade.Symbol] = append(session.Trades[trade.Symbol], trade)
})
// feed klines into the market data store
if environ.startTime == emptyTime {
environ.startTime = time.Now()
}
for symbol := range session.loadedSymbols {
2020-10-28 01:43:19 +00:00
marketDataStore, ok := session.marketDataStores[symbol]
if !ok {
return errors.Errorf("symbol %s is not defined", symbol)
2020-10-28 01:43:19 +00:00
}
var lastPriceTime time.Time
2020-10-28 01:43:19 +00:00
for interval := range types.SupportedIntervals {
kLines, err := session.Exchange.QueryKLines(ctx, symbol, interval, types.KLineQueryOptions{
EndTime: &environ.startTime,
Limit: 500, // indicators need at least 100
2020-10-28 01:43:19 +00:00
})
if err != nil {
return err
}
if len(kLines) == 0 {
log.Warnf("no kline data for interval %s", interval)
continue
}
// update last prices by the given kline
lastKLine := kLines[len(kLines) - 1]
if lastPriceTime == emptyTime {
session.lastPrices[symbol] = lastKLine.Close
lastPriceTime = lastKLine.EndTime
} else if lastPriceTime.Before(lastKLine.EndTime) {
session.lastPrices[symbol] = lastKLine.Close
lastPriceTime = lastKLine.EndTime
}
2020-10-28 01:43:19 +00:00
for _, k := range kLines {
// let market data store trigger the update, so that the indicator could be updated too.
marketDataStore.AddKLine(k)
}
}
2020-10-16 02:14:36 +00:00
}
if environ.TradeService != nil {
session.Stream.OnTradeUpdate(func(trade types.Trade) {
if err := environ.TradeService.Insert(trade); err != nil {
log.WithError(err).Errorf("trade insert error: %+v", trade)
}
})
}
// TODO: move market data store dispatch to here, use one callback to dispatch the market data
// session.Stream.OnKLineClosed(func(kline types.KLine) { })
}
return nil
}
2020-10-30 21:21:17 +00:00
// configure notification rules
// for symbol-based routes, we should register the same symbol rules for each session.
// for session-based routes, we should set the fixed callbacks for each session
func (environ *Environment) ConfigureNotification(conf *NotificationConfig) {
// configure routing here
if conf.SymbolChannels != nil {
environ.SymbolChannelRouter.AddRoute(conf.SymbolChannels)
}
if conf.SessionChannels != nil {
environ.SessionChannelRouter.AddRoute(conf.SessionChannels)
}
if conf.Routing != nil {
// configure passive object notification routing
switch conf.Routing.Trade {
case "$session":
defaultTradeUpdateHandler := func(trade types.Trade) {
text := util.Render(TemplateTradeReport, trade)
environ.Notify(text, &trade)
2020-10-30 21:21:17 +00:00
}
for name := range environ.sessions {
session := environ.sessions[name]
// if we can route session name to channel successfully...
channel, ok := environ.SessionChannelRouter.Route(name)
if ok {
session.Stream.OnTradeUpdate(func(trade types.Trade) {
text := util.Render(TemplateTradeReport, trade)
environ.NotifyTo(channel, text, &trade)
2020-10-30 21:21:17 +00:00
})
} else {
session.Stream.OnTradeUpdate(defaultTradeUpdateHandler)
}
}
case "$symbol":
// configure object routes for Trade
environ.ObjectChannelRouter.Route(func(obj interface{}) (channel string, ok bool) {
trade, matched := obj.(*types.Trade)
if !matched {
return
}
channel, ok = environ.SymbolChannelRouter.Route(trade.Symbol)
return
})
// use same handler for each session
handler := func(trade types.Trade) {
text := util.Render(TemplateTradeReport, trade)
channel, ok := environ.RouteObject(&trade)
2020-10-30 21:21:17 +00:00
if ok {
environ.NotifyTo(channel, text, &trade)
2020-10-30 21:21:17 +00:00
} else {
environ.Notify(text, &trade)
2020-10-30 21:21:17 +00:00
}
}
for _, session := range environ.sessions {
session.Stream.OnTradeUpdate(handler)
}
}
switch conf.Routing.Order {
case "$session":
defaultOrderUpdateHandler := func(order types.Order) {
text := util.Render(TemplateOrderReport, order)
environ.Notify(text, &order)
2020-10-30 21:21:17 +00:00
}
for name := range environ.sessions {
session := environ.sessions[name]
// if we can route session name to channel successfully...
channel, ok := environ.SessionChannelRouter.Route(name)
if ok {
session.Stream.OnOrderUpdate(func(order types.Order) {
text := util.Render(TemplateOrderReport, order)
environ.NotifyTo(channel, text, &order)
2020-10-30 21:21:17 +00:00
})
} else {
session.Stream.OnOrderUpdate(defaultOrderUpdateHandler)
}
}
case "$symbol":
// add object route
environ.ObjectChannelRouter.Route(func(obj interface{}) (channel string, ok bool) {
order, matched := obj.(*types.Order)
if !matched {
return
}
channel, ok = environ.SymbolChannelRouter.Route(order.Symbol)
return
})
// use same handler for each session
handler := func(order types.Order) {
text := util.Render(TemplateOrderReport, order)
channel, ok := environ.RouteObject(&order)
2020-10-30 21:21:17 +00:00
if ok {
environ.NotifyTo(channel, text, &order)
2020-10-30 21:21:17 +00:00
} else {
environ.Notify(text, &order)
2020-10-30 21:21:17 +00:00
}
}
for _, session := range environ.sessions {
session.Stream.OnOrderUpdate(handler)
}
}
switch conf.Routing.SubmitOrder {
case "$symbol":
// add object route
environ.ObjectChannelRouter.Route(func(obj interface{}) (channel string, ok bool) {
order, matched := obj.(*types.SubmitOrder)
if !matched {
return
}
channel, ok = environ.SymbolChannelRouter.Route(order.Symbol)
return
})
}
// currently not used
switch conf.Routing.PnL {
case "$symbol":
environ.ObjectChannelRouter.Route(func(obj interface{}) (channel string, ok bool) {
report, matched := obj.(*pnl.AverageCostPnlReport)
if !matched {
return
}
channel, ok = environ.SymbolChannelRouter.Route(report.Symbol)
return
})
}
}
}
// SyncTradesFrom overrides the default trade scan time (-7 days)
func (environ *Environment) SyncTradesFrom(t time.Time) *Environment {
environ.tradeScanTime = t
return environ
}
func (environ *Environment) Connect(ctx context.Context) error {
for n := range environ.sessions {
// avoid using the placeholder variable for the session because we use that in the callbacks
var session = environ.sessions[n]
var logger = log.WithField("session", n)
2020-10-19 14:26:43 +00:00
if len(session.Subscriptions) == 0 {
logger.Warnf("no subscriptions, exchange session %s will not be connected", session.Name)
2020-10-19 14:26:43 +00:00
continue
}
// add the subscribe requests to the stream
for _, s := range session.Subscriptions {
logger.Infof("subscribing %s %s %v", s.Symbol, s.Channel, s.Options)
session.Stream.Subscribe(s.Channel, s.Symbol, s.Options)
}
logger.Infof("connecting session %s...", session.Name)
2020-10-16 02:14:36 +00:00
if err := session.Stream.Connect(ctx); err != nil {
return err
}
}
return nil
}