19 changed files with 346 additions and 41 deletions
@ -0,0 +1,101 @@
|
||||
package main |
||||
|
||||
import ( |
||||
"fmt" |
||||
"net" |
||||
"sig-pub/api/pb" |
||||
"sig-pub/internal/trading" |
||||
"sig-pub/pkg/client" |
||||
"sig-pub/pkg/config" |
||||
"sig-pub/pkg/grpc/discovery" |
||||
"sig-pub/pkg/grpc/interceptor" |
||||
"sig-pub/pkg/utils/exit" |
||||
"sig-pub/pkg/zlog" |
||||
|
||||
"github.com/hashicorp/consul/api" |
||||
_ "github.com/mostynb/go-grpc-compression/snappy" // 注册grpc snappy compress
|
||||
"google.golang.org/grpc" |
||||
"google.golang.org/grpc/credentials/insecure" |
||||
"google.golang.org/grpc/reflection" |
||||
) |
||||
|
||||
type TradingConf struct { |
||||
Register discovery.Server |
||||
GrpcReflection bool |
||||
} |
||||
|
||||
func main() { |
||||
// load config
|
||||
conf := config.MustLoadConfig(new(config.Configuration), "config/config.toml") |
||||
tradingConf := config.MustLoadConfig(new(TradingConf), "config/trading.toml") |
||||
|
||||
// consul 配置
|
||||
cc := api.DefaultConfig() |
||||
cc.Address = conf.Consul.Address |
||||
consulClient, err := api.NewClient(cc) |
||||
if err != nil { |
||||
panic(fmt.Errorf("consul client error: %v", err)) |
||||
} |
||||
dis := discovery.NewConsulDiscovery(consulClient) |
||||
resolver := dis.Resolver() |
||||
|
||||
// new market grpc client
|
||||
marketClient, err := client.NewMarketClient( |
||||
grpc.WithResolvers(resolver), |
||||
grpc.WithTransportCredentials(insecure.NewCredentials())) |
||||
if err != nil { |
||||
panic(err) |
||||
} |
||||
tradeInstanceAside := client.NewTradeInstanceAside(marketClient) |
||||
|
||||
// new exchange grpc client
|
||||
exchangeClient, err := client.NewExchangeClient( |
||||
grpc.WithResolvers(resolver), |
||||
grpc.WithTransportCredentials(insecure.NewCredentials())) |
||||
if err != nil { |
||||
panic(err) |
||||
} |
||||
|
||||
// services
|
||||
tradingService := trading.NewTradingService(tradeInstanceAside, exchangeClient) |
||||
if err := tradingService.Init(); err != nil { |
||||
panic(err) |
||||
} |
||||
|
||||
// grpc server
|
||||
tradingGrpcServer := trading.NewTradingGrpcServer() |
||||
if err := tradingGrpcServer.Init(); err != nil { |
||||
panic(err) |
||||
} |
||||
grpcServer := grpc.NewServer(config.GetGrpcOptions( |
||||
conf.Grpc, |
||||
grpc.UnaryInterceptor(interceptor.RecoverInterceptor), |
||||
)...) |
||||
if tradingConf.GrpcReflection { |
||||
// 注册反射服务
|
||||
reflection.Register(grpcServer) |
||||
} |
||||
pb.RegisterTradingServiceServer(grpcServer, tradingGrpcServer) |
||||
exit.AddHook(grpcServer.GracefulStop, exit.WithOrderFront()) |
||||
|
||||
// consul 服务注册
|
||||
register := tradingConf.Register |
||||
register.Name = pb.TradingService_ServiceDesc.ServiceName |
||||
if err := dis.Registry(grpcServer, register); err != nil { |
||||
panic(err) |
||||
} |
||||
|
||||
// run grpc server
|
||||
go func() { |
||||
listen, err := net.Listen("tcp", tradingConf.Register.Addr) |
||||
if err != nil { |
||||
panic(err) |
||||
} |
||||
zlog.Infof("%s grpc server running %s\n", register.Name, listen.Addr().String()) |
||||
if err := grpcServer.Serve(listen); err != nil { |
||||
panic(err) |
||||
} |
||||
}() |
||||
|
||||
exit.Await() |
||||
} |
||||
@ -0,0 +1,7 @@
|
||||
|
||||
grpcReflection = true # 注册grpc反射服务 |
||||
|
||||
[register] |
||||
nodeId = 1 # grpc服务节点id, 多实例唯一 |
||||
addr = ":8031" # grpc 服务端口 |
||||
attrs = { weight = 10 } # grpc 服务权重 |
||||
@ -0,0 +1,17 @@
|
||||
package trading |
||||
|
||||
import ( |
||||
"sig-pub/api/pb" |
||||
) |
||||
|
||||
type TradingGrpcServer struct { |
||||
pb.UnimplementedTradingServiceServer |
||||
} |
||||
|
||||
func NewTradingGrpcServer() *TradingGrpcServer { |
||||
return &TradingGrpcServer{} |
||||
} |
||||
|
||||
func (svr TradingGrpcServer) Init() (err error) { |
||||
return |
||||
} |
||||
@ -0,0 +1,136 @@
|
||||
package trading |
||||
|
||||
import ( |
||||
"context" |
||||
"io" |
||||
"sig-pub/api/pb" |
||||
"sig-pub/pkg/client" |
||||
"sig-pub/pkg/types" |
||||
"sig-pub/pkg/zlog" |
||||
"sync" |
||||
"time" |
||||
|
||||
"google.golang.org/grpc" |
||||
) |
||||
|
||||
type TradingService struct { |
||||
marketClientAside *client.TradeInstanceAside |
||||
exchangeClient pb.ExchangeServiceClient |
||||
|
||||
klineStore *KlineStore |
||||
|
||||
subKlineLock sync.Mutex |
||||
subKlineInsts [3][]string |
||||
subKlineStream grpc.BidiStreamingClient[pb.ReqStreamSubscribeKline, pb.RspStreamSubscribeKline] |
||||
} |
||||
|
||||
func NewTradingService( |
||||
marketClientAside *client.TradeInstanceAside, |
||||
exchangeClient pb.ExchangeServiceClient, |
||||
) *TradingService { |
||||
return &TradingService{ |
||||
marketClientAside: marketClientAside, |
||||
exchangeClient: exchangeClient, |
||||
} |
||||
} |
||||
|
||||
// 初始化历史k线, 订阅实时k线
|
||||
func (svr *TradingService) Init() (err error) { |
||||
// 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.subscribeKlines(false) |
||||
return |
||||
} |
||||
|
||||
func (svr *TradingService) subscribeKlines(reconnect bool) { |
||||
defer func() { |
||||
svr.subKlineStream = nil |
||||
go svr.subscribeKlines(true) |
||||
}() |
||||
|
||||
if reconnect { |
||||
zlog.Infof("subscribeKlines will reconnect after 5s") |
||||
time.Sleep(5 * time.Second) |
||||
} |
||||
|
||||
// svr.subKlineLock.Lock()
|
||||
// defer svr.subKlineLock.Unlock()
|
||||
|
||||
stream, err := svr.exchangeClient.SubscribeKline(context.Background()) |
||||
if err != nil { |
||||
zlog.Error("subscribeKlines reqeust error: ", err) |
||||
return |
||||
} |
||||
|
||||
if svr.subKlineStream != nil { |
||||
svr.subKlineStream.CloseSend() |
||||
} |
||||
|
||||
// 发送订阅消息
|
||||
exchanges := []pb.ExchangeType{pb.ExchangeType_OKX} |
||||
for _, exchange := range exchanges { |
||||
instIds := svr.subKlineInsts[exchange] |
||||
|
||||
msg := &pb.ReqStreamSubscribeKline{ |
||||
SubType: pb.SubscribeType_Subscribe, |
||||
Exchanges: []pb.ExchangeType{exchange}, |
||||
InstIds: instIds, |
||||
Intervals: 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) |
||||
|
||||
svr.klineStore.Update(msg.Kline.Exchange, msg.Kline.InstId, kline) |
||||
} |
||||
} |
||||
} |
||||
|
||||
var SubscribeKlineIntervals = []string{ |
||||
string(types.Interval1s), |
||||
string(types.Interval1m), |
||||
string(types.Interval3m), |
||||
string(types.Interval5m), |
||||
string(types.Interval15m), |
||||
string(types.Interval30m), |
||||
string(types.Interval1h), |
||||
string(types.Interval2h), |
||||
string(types.Interval4h), |
||||
string(types.Interval6h), |
||||
string(types.Interval12h), |
||||
string(types.Interval1d), |
||||
string(types.Interval1d), |
||||
string(types.Interval2d), |
||||
string(types.Interval3d), |
||||
string(types.Interval5d), |
||||
string(types.Interval1w), |
||||
} |
||||
@ -0,0 +1,30 @@
|
||||
package client |
||||
|
||||
import ( |
||||
"sig-pub/api/pb" |
||||
"sig-pub/pkg/grpc/discovery" |
||||
|
||||
"google.golang.org/grpc" |
||||
) |
||||
|
||||
// NewMarketClient new MarketServer grpc client
|
||||
func NewMarketClient(options ...grpc.DialOption) (marketClient pb.MarketServiceClient, err error) { |
||||
grpcUrl := discovery.ConsulDialUrl(pb.MarketService_ServiceDesc.ServiceName) |
||||
grpcConn, err := grpc.NewClient(grpcUrl, options...) |
||||
if err != nil { |
||||
panic(err) |
||||
} |
||||
marketClient = pb.NewMarketServiceClient(grpcConn) |
||||
return |
||||
} |
||||
|
||||
// NewExchangeClient new ExchangeServer grpc client
|
||||
func NewExchangeClient(options ...grpc.DialOption) (exchangeClient pb.ExchangeServiceClient, err error) { |
||||
grpcUrl := discovery.ConsulDialUrl(pb.ExchangeService_ServiceDesc.ServiceName) |
||||
grpcConn, err := grpc.NewClient(grpcUrl, options...) |
||||
if err != nil { |
||||
panic(err) |
||||
} |
||||
exchangeClient = pb.NewExchangeServiceClient(grpcConn) |
||||
return |
||||
} |
||||
@ -1,4 +1,4 @@
|
||||
package aside |
||||
package client |
||||
|
||||
import ( |
||||
"context" |
||||
Loading…
Reference in new issue