bbgo_origin/pkg/grpc/server.go

364 lines
9.0 KiB
Go
Raw Permalink Normal View History

2022-04-08 11:21:57 +00:00
package grpc
import (
"context"
2022-04-12 07:53:40 +00:00
"fmt"
2022-04-08 11:21:57 +00:00
"net"
2022-04-15 07:49:24 +00:00
"strconv"
"time"
2022-04-08 11:21:57 +00:00
"github.com/pkg/errors"
2022-04-12 08:55:42 +00:00
log "github.com/sirupsen/logrus"
"go.uber.org/multierr"
2022-04-08 11:21:57 +00:00
"google.golang.org/grpc"
"google.golang.org/grpc/reflection"
"github.com/c9s/bbgo/pkg/bbgo"
"github.com/c9s/bbgo/pkg/pb"
"github.com/c9s/bbgo/pkg/types"
)
2022-04-15 06:53:50 +00:00
type TradingService struct {
Config *bbgo.Config
Environ *bbgo.Environment
Trader *bbgo.Trader
pb.UnimplementedTradingServiceServer
}
func (s *TradingService) SubmitOrder(ctx context.Context, request *pb.SubmitOrderRequest) (*pb.SubmitOrderResponse, error) {
sessionName := request.Session
if len(sessionName) == 0 {
return nil, fmt.Errorf("session name can not be empty")
}
session, ok := s.Environ.Session(sessionName)
if !ok {
return nil, fmt.Errorf("session %s not found", sessionName)
}
submitOrders := toSubmitOrders(request.SubmitOrders)
for i := range submitOrders {
if market, ok := session.Market(submitOrders[i].Symbol); ok {
submitOrders[i].Market = market
} else {
log.Warnf("session %s market %s not found", sessionName, submitOrders[i].Symbol)
}
}
createdOrders, errIdx, err := bbgo.BatchPlaceOrder(ctx, session.Exchange, submitOrders...)
if len(errIdx) > 0 {
createdOrders2, err2 := bbgo.BatchRetryPlaceOrder(ctx, session.Exchange, errIdx, submitOrders...)
if err2 != nil {
err = multierr.Append(err, err2)
} else {
createdOrders = append(createdOrders, createdOrders2...)
}
2022-04-15 06:53:50 +00:00
}
// convert response
2022-04-15 07:03:00 +00:00
resp := &pb.SubmitOrderResponse{
Session: sessionName,
Orders: nil,
}
2022-04-15 07:03:00 +00:00
for _, createdOrder := range createdOrders {
resp.Orders = append(resp.Orders, transOrder(session, createdOrder))
}
return resp, err
2022-04-15 06:53:50 +00:00
}
func (s *TradingService) CancelOrder(ctx context.Context, request *pb.CancelOrderRequest) (*pb.CancelOrderResponse, error) {
2022-04-15 07:49:24 +00:00
sessionName := request.Session
if len(sessionName) == 0 {
return nil, fmt.Errorf("session name can not be empty")
}
session, ok := s.Environ.Session(sessionName)
if !ok {
return nil, fmt.Errorf("session %s not found", sessionName)
}
uuidOrderID := ""
orderID, err := strconv.ParseUint(request.OrderId, 10, 64)
if err != nil {
// TODO: validate uuid
uuidOrderID = request.OrderId
}
session.Exchange.CancelOrders(ctx, types.Order{
SubmitOrder: types.SubmitOrder{
ClientOrderID: request.ClientOrderId,
},
OrderID: orderID,
UUID: uuidOrderID,
})
resp := &pb.CancelOrderResponse{}
return resp, nil
2022-04-15 06:53:50 +00:00
}
func (s *TradingService) QueryOrder(ctx context.Context, request *pb.QueryOrderRequest) (*pb.QueryOrderResponse, error) {
panic("implement me")
}
func (s *TradingService) QueryOrders(ctx context.Context, request *pb.QueryOrdersRequest) (*pb.QueryOrdersResponse, error) {
panic("implement me")
}
func (s *TradingService) QueryTrades(ctx context.Context, request *pb.QueryTradesRequest) (*pb.QueryTradesResponse, error) {
panic("implement me")
}
2022-04-13 07:29:23 +00:00
type UserDataService struct {
2022-04-08 11:21:57 +00:00
Config *bbgo.Config
Environ *bbgo.Environment
Trader *bbgo.Trader
pb.UnimplementedUserDataServiceServer
}
2022-04-13 07:29:23 +00:00
func (s *UserDataService) Subscribe(request *pb.UserDataRequest, server pb.UserDataService_SubscribeServer) error {
2022-04-14 16:13:19 +00:00
sessionName := request.Session
if len(sessionName) == 0 {
return fmt.Errorf("session name can not be empty")
}
2022-04-14 16:13:19 +00:00
session, ok := s.Environ.Session(sessionName)
if !ok {
return fmt.Errorf("session %s not found", sessionName)
2022-04-13 07:29:23 +00:00
}
2022-04-14 16:13:19 +00:00
userDataStream := session.Exchange.NewStream()
userDataStream.OnOrderUpdate(func(order types.Order) {
err := server.Send(&pb.UserData{
Channel: pb.Channel_ORDER,
Event: pb.Event_UPDATE,
Orders: []*pb.Order{transOrder(session, order)},
})
if err != nil {
log.WithError(err).Errorf("grpc: can not send user data")
}
})
userDataStream.OnTradeUpdate(func(trade types.Trade) {
err := server.Send(&pb.UserData{
Channel: pb.Channel_TRADE,
Event: pb.Event_UPDATE,
Trades: []*pb.Trade{transTrade(session, trade)},
})
if err != nil {
log.WithError(err).Errorf("grpc: can not send user data")
}
})
balanceHandler := func(balances types.BalanceMap) {
err := server.Send(&pb.UserData{
Channel: pb.Channel_BALANCE,
Event: pb.Event_UPDATE,
Balances: transBalances(session, balances),
})
if err != nil {
log.WithError(err).Errorf("grpc: can not send user data")
}
}
userDataStream.OnBalanceUpdate(balanceHandler)
userDataStream.OnBalanceSnapshot(balanceHandler)
ctx := server.Context()
balances, err := session.Exchange.QueryAccountBalances(ctx)
if err != nil {
return err
}
err = server.Send(&pb.UserData{
Channel: pb.Channel_BALANCE,
Event: pb.Event_SNAPSHOT,
Balances: transBalances(session, balances),
})
if err != nil {
log.WithError(err).Errorf("grpc: can not send user data")
}
go userDataStream.Connect(ctx)
defer func() {
2022-04-14 16:13:19 +00:00
if err := userDataStream.Close(); err != nil {
log.WithError(err).Errorf("user data stream close error")
}
}()
<-ctx.Done()
return nil
2022-04-08 11:21:57 +00:00
}
2022-04-13 07:29:23 +00:00
type MarketDataService struct {
Config *bbgo.Config
Environ *bbgo.Environment
Trader *bbgo.Trader
pb.UnimplementedMarketDataServiceServer
}
func (s *MarketDataService) Subscribe(request *pb.SubscribeRequest, server pb.MarketDataService_SubscribeServer) error {
2022-04-12 07:53:40 +00:00
exchangeSubscriptions := map[string][]types.Subscription{}
for _, sub := range request.Subscriptions {
session, ok := s.Environ.Session(sub.Exchange)
if !ok {
return fmt.Errorf("exchange %s not found", sub.Exchange)
}
2022-04-13 05:06:26 +00:00
ss, err := toSubscriptions(sub)
if err != nil {
return err
2022-04-12 07:53:40 +00:00
}
2022-04-13 05:06:26 +00:00
exchangeSubscriptions[session.Name] = append(exchangeSubscriptions[session.Name], ss)
2022-04-12 07:53:40 +00:00
}
2022-04-12 09:48:30 +00:00
streamPool := map[string]types.Stream{}
2022-04-12 07:53:40 +00:00
for sessionName, subs := range exchangeSubscriptions {
2022-04-12 15:25:41 +00:00
session, ok := s.Environ.Session(sessionName)
if !ok {
log.Errorf("session %s not found", sessionName)
continue
}
2022-04-12 08:55:42 +00:00
2022-04-12 15:25:41 +00:00
stream := session.Exchange.NewStream()
stream.SetPublicOnly()
for _, sub := range subs {
log.Infof("%s subscribe %s %s %+v", sessionName, sub.Channel, sub.Symbol, sub.Options)
stream.Subscribe(sub.Channel, sub.Symbol, sub.Options)
2022-04-12 07:53:40 +00:00
}
2022-04-12 15:25:41 +00:00
2022-04-13 03:52:51 +00:00
stream.OnMarketTrade(func(trade types.Trade) {
if err := server.Send(transMarketTrade(session, trade)); err != nil {
log.WithError(err).Error("grpc stream send error")
}
})
2022-04-12 15:25:41 +00:00
stream.OnBookSnapshot(func(book types.SliceOrderBook) {
if err := server.Send(transBook(session, book, pb.Event_SNAPSHOT)); err != nil {
log.WithError(err).Error("grpc stream send error")
}
})
stream.OnBookUpdate(func(book types.SliceOrderBook) {
if err := server.Send(transBook(session, book, pb.Event_UPDATE)); err != nil {
log.WithError(err).Error("grpc stream send error")
}
})
stream.OnKLineClosed(func(kline types.KLine) {
err := server.Send(transKLineResponse(session, kline))
2022-04-12 15:25:41 +00:00
if err != nil {
log.WithError(err).Error("grpc stream send error")
}
})
streamPool[sessionName] = stream
2022-04-12 09:48:30 +00:00
}
2022-04-13 03:52:51 +00:00
for _, stream := range streamPool {
2022-04-12 09:48:30 +00:00
go stream.Connect(server.Context())
2022-04-12 07:53:40 +00:00
}
2022-04-13 03:52:51 +00:00
defer func() {
for _, stream := range streamPool {
if err := stream.Close(); err != nil {
log.WithError(err).Errorf("market data stream close error")
2022-04-13 03:52:51 +00:00
}
}
}()
ctx := server.Context()
<-ctx.Done()
return ctx.Err()
2022-04-08 11:21:57 +00:00
}
2022-04-13 07:29:23 +00:00
func (s *MarketDataService) QueryKLines(ctx context.Context, request *pb.QueryKLinesRequest) (*pb.QueryKLinesResponse, error) {
2022-04-08 11:21:57 +00:00
exchangeName, err := types.ValidExchangeName(request.Exchange)
if err != nil {
return nil, err
}
for _, session := range s.Environ.Sessions() {
if session.ExchangeName == exchangeName {
response := &pb.QueryKLinesResponse{
Klines: nil,
Error: nil,
}
options := types.KLineQueryOptions{
Limit: int(request.Limit),
}
endTime := time.Now()
if request.EndTime != 0 {
endTime = time.Unix(request.EndTime, 0)
}
options.EndTime = &endTime
if request.StartTime != 0 {
startTime := time.Unix(request.StartTime, 0)
options.StartTime = &startTime
}
2022-04-08 11:21:57 +00:00
klines, err := session.Exchange.QueryKLines(ctx, request.Symbol, types.Interval(request.Interval), options)
if err != nil {
return nil, err
}
for _, kline := range klines {
response.Klines = append(response.Klines, transKLine(session, kline))
2022-04-08 11:21:57 +00:00
}
return response, nil
}
}
return nil, nil
}
2022-04-13 07:29:23 +00:00
type Server struct {
Config *bbgo.Config
Environ *bbgo.Environment
Trader *bbgo.Trader
}
2022-04-08 11:21:57 +00:00
func (s *Server) ListenAndServe(bind string) error {
conn, err := net.Listen("tcp", bind)
if err != nil {
return errors.Wrapf(err, "failed to bind network at %s", bind)
}
var grpcServer = grpc.NewServer()
2022-04-13 07:29:23 +00:00
pb.RegisterMarketDataServiceServer(grpcServer, &MarketDataService{
Config: s.Config,
Environ: s.Environ,
Trader: s.Trader,
})
2022-04-16 15:57:53 +00:00
pb.RegisterTradingServiceServer(grpcServer, &TradingService{
Config: s.Config,
Environ: s.Environ,
Trader: s.Trader,
})
2022-04-13 07:29:23 +00:00
pb.RegisterUserDataServiceServer(grpcServer, &UserDataService{
Config: s.Config,
Environ: s.Environ,
Trader: s.Trader,
})
2022-04-08 11:21:57 +00:00
reflection.Register(grpcServer)
if err := grpcServer.Serve(conn); err != nil {
return errors.Wrap(err, "failed to serve grpc connections")
}
return nil
}