pb: add grpc files

This commit is contained in:
c9s 2022-04-07 19:28:05 +08:00
parent acaa68c049
commit 0885b1635a
4 changed files with 714 additions and 137 deletions

View File

@ -242,9 +242,11 @@ PROTOS := \
GRPC_GO_DEPS := $(subst .proto,.pb.go,$(PROTOS)) GRPC_GO_DEPS := $(subst .proto,.pb.go,$(PROTOS))
%.pb.go: %.proto .FORCE %.pb.go: %.proto .FORCE
protoc --go-grpc_out=. --go_out=. --proto_path=. $< protoc --go-grpc_out=. --go-grpc_opt=paths=source_relative --go_out=. --proto_path=. $<
grpc: $(GRPC_GO_DEPS) grpc-py grpc-go: $(GRPC_GO_DEPS)
grpc: grpc-go grpc-py
install-grpc-tools: install-grpc-tools:
go install google.golang.org/protobuf/cmd/protoc-gen-go@v1.26 go install google.golang.org/protobuf/cmd/protoc-gen-go@v1.26
@ -257,4 +259,4 @@ grpc-py:
--grpc_python_out=$(PWD)/python/bbgo \ --grpc_python_out=$(PWD)/python/bbgo \
$(PWD)/pkg/pb/bbgo.proto $(PWD)/pkg/pb/bbgo.proto
.PHONY: bbgo bbgo-slim-darwin bbgo-slim-darwin-amd64 bbgo-slim-darwin-arm64 bbgo-darwin version dist pack migrations static embed desktop .FORCE .PHONY: bbgo bbgo-slim-darwin bbgo-slim-darwin-amd64 bbgo-slim-darwin-arm64 bbgo-darwin version dist pack migrations static embed desktop grpc grpc-go grpc-py .FORCE

512
pkg/pb/bbgo_grpc.pb.go Normal file
View File

@ -0,0 +1,512 @@
// Code generated by protoc-gen-go-grpc. DO NOT EDIT.
package pb
import (
context "context"
grpc "google.golang.org/grpc"
codes "google.golang.org/grpc/codes"
status "google.golang.org/grpc/status"
)
// This is a compile-time assertion to ensure that this generated file
// is compatible with the grpc package it is being compiled against.
// Requires gRPC-Go v1.32.0 or later.
const _ = grpc.SupportPackageIsVersion7
// MarketDataServiceClient is the client API for MarketDataService service.
//
// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
type MarketDataServiceClient interface {
Subscribe(ctx context.Context, in *SubscribeRequest, opts ...grpc.CallOption) (MarketDataService_SubscribeClient, error)
QueryKLines(ctx context.Context, in *QueryKLinesRequest, opts ...grpc.CallOption) (*QueryKLinesResponse, error)
}
type marketDataServiceClient struct {
cc grpc.ClientConnInterface
}
func NewMarketDataServiceClient(cc grpc.ClientConnInterface) MarketDataServiceClient {
return &marketDataServiceClient{cc}
}
func (c *marketDataServiceClient) Subscribe(ctx context.Context, in *SubscribeRequest, opts ...grpc.CallOption) (MarketDataService_SubscribeClient, error) {
stream, err := c.cc.NewStream(ctx, &MarketDataService_ServiceDesc.Streams[0], "/pb.MarketDataService/Subscribe", opts...)
if err != nil {
return nil, err
}
x := &marketDataServiceSubscribeClient{stream}
if err := x.ClientStream.SendMsg(in); err != nil {
return nil, err
}
if err := x.ClientStream.CloseSend(); err != nil {
return nil, err
}
return x, nil
}
type MarketDataService_SubscribeClient interface {
Recv() (*SubscribeResponse, error)
grpc.ClientStream
}
type marketDataServiceSubscribeClient struct {
grpc.ClientStream
}
func (x *marketDataServiceSubscribeClient) Recv() (*SubscribeResponse, error) {
m := new(SubscribeResponse)
if err := x.ClientStream.RecvMsg(m); err != nil {
return nil, err
}
return m, nil
}
func (c *marketDataServiceClient) QueryKLines(ctx context.Context, in *QueryKLinesRequest, opts ...grpc.CallOption) (*QueryKLinesResponse, error) {
out := new(QueryKLinesResponse)
err := c.cc.Invoke(ctx, "/pb.MarketDataService/QueryKLines", in, out, opts...)
if err != nil {
return nil, err
}
return out, nil
}
// MarketDataServiceServer is the server API for MarketDataService service.
// All implementations must embed UnimplementedMarketDataServiceServer
// for forward compatibility
type MarketDataServiceServer interface {
Subscribe(*SubscribeRequest, MarketDataService_SubscribeServer) error
QueryKLines(context.Context, *QueryKLinesRequest) (*QueryKLinesResponse, error)
mustEmbedUnimplementedMarketDataServiceServer()
}
// UnimplementedMarketDataServiceServer must be embedded to have forward compatible implementations.
type UnimplementedMarketDataServiceServer struct {
}
func (UnimplementedMarketDataServiceServer) Subscribe(*SubscribeRequest, MarketDataService_SubscribeServer) error {
return status.Errorf(codes.Unimplemented, "method Subscribe not implemented")
}
func (UnimplementedMarketDataServiceServer) QueryKLines(context.Context, *QueryKLinesRequest) (*QueryKLinesResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method QueryKLines not implemented")
}
func (UnimplementedMarketDataServiceServer) mustEmbedUnimplementedMarketDataServiceServer() {}
// UnsafeMarketDataServiceServer may be embedded to opt out of forward compatibility for this service.
// Use of this interface is not recommended, as added methods to MarketDataServiceServer will
// result in compilation errors.
type UnsafeMarketDataServiceServer interface {
mustEmbedUnimplementedMarketDataServiceServer()
}
func RegisterMarketDataServiceServer(s grpc.ServiceRegistrar, srv MarketDataServiceServer) {
s.RegisterService(&MarketDataService_ServiceDesc, srv)
}
func _MarketDataService_Subscribe_Handler(srv interface{}, stream grpc.ServerStream) error {
m := new(SubscribeRequest)
if err := stream.RecvMsg(m); err != nil {
return err
}
return srv.(MarketDataServiceServer).Subscribe(m, &marketDataServiceSubscribeServer{stream})
}
type MarketDataService_SubscribeServer interface {
Send(*SubscribeResponse) error
grpc.ServerStream
}
type marketDataServiceSubscribeServer struct {
grpc.ServerStream
}
func (x *marketDataServiceSubscribeServer) Send(m *SubscribeResponse) error {
return x.ServerStream.SendMsg(m)
}
func _MarketDataService_QueryKLines_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(QueryKLinesRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(MarketDataServiceServer).QueryKLines(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: "/pb.MarketDataService/QueryKLines",
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(MarketDataServiceServer).QueryKLines(ctx, req.(*QueryKLinesRequest))
}
return interceptor(ctx, in, info, handler)
}
// MarketDataService_ServiceDesc is the grpc.ServiceDesc for MarketDataService service.
// It's only intended for direct use with grpc.RegisterService,
// and not to be introspected or modified (even as a copy)
var MarketDataService_ServiceDesc = grpc.ServiceDesc{
ServiceName: "pb.MarketDataService",
HandlerType: (*MarketDataServiceServer)(nil),
Methods: []grpc.MethodDesc{
{
MethodName: "QueryKLines",
Handler: _MarketDataService_QueryKLines_Handler,
},
},
Streams: []grpc.StreamDesc{
{
StreamName: "Subscribe",
Handler: _MarketDataService_Subscribe_Handler,
ServerStreams: true,
},
},
Metadata: "pkg/pb/bbgo.proto",
}
// UserDataServiceClient is the client API for UserDataService service.
//
// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
type UserDataServiceClient interface {
// should support streaming
SubscribeUserData(ctx context.Context, in *Empty, opts ...grpc.CallOption) (UserDataService_SubscribeUserDataClient, error)
}
type userDataServiceClient struct {
cc grpc.ClientConnInterface
}
func NewUserDataServiceClient(cc grpc.ClientConnInterface) UserDataServiceClient {
return &userDataServiceClient{cc}
}
func (c *userDataServiceClient) SubscribeUserData(ctx context.Context, in *Empty, opts ...grpc.CallOption) (UserDataService_SubscribeUserDataClient, error) {
stream, err := c.cc.NewStream(ctx, &UserDataService_ServiceDesc.Streams[0], "/pb.UserDataService/SubscribeUserData", opts...)
if err != nil {
return nil, err
}
x := &userDataServiceSubscribeUserDataClient{stream}
if err := x.ClientStream.SendMsg(in); err != nil {
return nil, err
}
if err := x.ClientStream.CloseSend(); err != nil {
return nil, err
}
return x, nil
}
type UserDataService_SubscribeUserDataClient interface {
Recv() (*SubscribeResponse, error)
grpc.ClientStream
}
type userDataServiceSubscribeUserDataClient struct {
grpc.ClientStream
}
func (x *userDataServiceSubscribeUserDataClient) Recv() (*SubscribeResponse, error) {
m := new(SubscribeResponse)
if err := x.ClientStream.RecvMsg(m); err != nil {
return nil, err
}
return m, nil
}
// UserDataServiceServer is the server API for UserDataService service.
// All implementations must embed UnimplementedUserDataServiceServer
// for forward compatibility
type UserDataServiceServer interface {
// should support streaming
SubscribeUserData(*Empty, UserDataService_SubscribeUserDataServer) error
mustEmbedUnimplementedUserDataServiceServer()
}
// UnimplementedUserDataServiceServer must be embedded to have forward compatible implementations.
type UnimplementedUserDataServiceServer struct {
}
func (UnimplementedUserDataServiceServer) SubscribeUserData(*Empty, UserDataService_SubscribeUserDataServer) error {
return status.Errorf(codes.Unimplemented, "method SubscribeUserData not implemented")
}
func (UnimplementedUserDataServiceServer) mustEmbedUnimplementedUserDataServiceServer() {}
// UnsafeUserDataServiceServer may be embedded to opt out of forward compatibility for this service.
// Use of this interface is not recommended, as added methods to UserDataServiceServer will
// result in compilation errors.
type UnsafeUserDataServiceServer interface {
mustEmbedUnimplementedUserDataServiceServer()
}
func RegisterUserDataServiceServer(s grpc.ServiceRegistrar, srv UserDataServiceServer) {
s.RegisterService(&UserDataService_ServiceDesc, srv)
}
func _UserDataService_SubscribeUserData_Handler(srv interface{}, stream grpc.ServerStream) error {
m := new(Empty)
if err := stream.RecvMsg(m); err != nil {
return err
}
return srv.(UserDataServiceServer).SubscribeUserData(m, &userDataServiceSubscribeUserDataServer{stream})
}
type UserDataService_SubscribeUserDataServer interface {
Send(*SubscribeResponse) error
grpc.ServerStream
}
type userDataServiceSubscribeUserDataServer struct {
grpc.ServerStream
}
func (x *userDataServiceSubscribeUserDataServer) Send(m *SubscribeResponse) error {
return x.ServerStream.SendMsg(m)
}
// UserDataService_ServiceDesc is the grpc.ServiceDesc for UserDataService service.
// It's only intended for direct use with grpc.RegisterService,
// and not to be introspected or modified (even as a copy)
var UserDataService_ServiceDesc = grpc.ServiceDesc{
ServiceName: "pb.UserDataService",
HandlerType: (*UserDataServiceServer)(nil),
Methods: []grpc.MethodDesc{},
Streams: []grpc.StreamDesc{
{
StreamName: "SubscribeUserData",
Handler: _UserDataService_SubscribeUserData_Handler,
ServerStreams: true,
},
},
Metadata: "pkg/pb/bbgo.proto",
}
// TradingServiceClient is the client API for TradingService service.
//
// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
type TradingServiceClient interface {
// request-response
SubmitOrder(ctx context.Context, in *SubmitOrderRequest, opts ...grpc.CallOption) (*SubmitOrderResponse, error)
CancelOrder(ctx context.Context, in *CancelOrderRequest, opts ...grpc.CallOption) (*CancelOrderResponse, error)
QueryOrder(ctx context.Context, in *QueryOrderRequest, opts ...grpc.CallOption) (*QueryOrderResponse, error)
QueryOrders(ctx context.Context, in *QueryOrdersRequest, opts ...grpc.CallOption) (*QueryOrdersResponse, error)
QueryTrades(ctx context.Context, in *QueryTradesRequest, opts ...grpc.CallOption) (*QueryTradesResponse, error)
}
type tradingServiceClient struct {
cc grpc.ClientConnInterface
}
func NewTradingServiceClient(cc grpc.ClientConnInterface) TradingServiceClient {
return &tradingServiceClient{cc}
}
func (c *tradingServiceClient) SubmitOrder(ctx context.Context, in *SubmitOrderRequest, opts ...grpc.CallOption) (*SubmitOrderResponse, error) {
out := new(SubmitOrderResponse)
err := c.cc.Invoke(ctx, "/pb.TradingService/SubmitOrder", in, out, opts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *tradingServiceClient) CancelOrder(ctx context.Context, in *CancelOrderRequest, opts ...grpc.CallOption) (*CancelOrderResponse, error) {
out := new(CancelOrderResponse)
err := c.cc.Invoke(ctx, "/pb.TradingService/CancelOrder", in, out, opts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *tradingServiceClient) QueryOrder(ctx context.Context, in *QueryOrderRequest, opts ...grpc.CallOption) (*QueryOrderResponse, error) {
out := new(QueryOrderResponse)
err := c.cc.Invoke(ctx, "/pb.TradingService/QueryOrder", in, out, opts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *tradingServiceClient) QueryOrders(ctx context.Context, in *QueryOrdersRequest, opts ...grpc.CallOption) (*QueryOrdersResponse, error) {
out := new(QueryOrdersResponse)
err := c.cc.Invoke(ctx, "/pb.TradingService/QueryOrders", in, out, opts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *tradingServiceClient) QueryTrades(ctx context.Context, in *QueryTradesRequest, opts ...grpc.CallOption) (*QueryTradesResponse, error) {
out := new(QueryTradesResponse)
err := c.cc.Invoke(ctx, "/pb.TradingService/QueryTrades", in, out, opts...)
if err != nil {
return nil, err
}
return out, nil
}
// TradingServiceServer is the server API for TradingService service.
// All implementations must embed UnimplementedTradingServiceServer
// for forward compatibility
type TradingServiceServer interface {
// request-response
SubmitOrder(context.Context, *SubmitOrderRequest) (*SubmitOrderResponse, error)
CancelOrder(context.Context, *CancelOrderRequest) (*CancelOrderResponse, error)
QueryOrder(context.Context, *QueryOrderRequest) (*QueryOrderResponse, error)
QueryOrders(context.Context, *QueryOrdersRequest) (*QueryOrdersResponse, error)
QueryTrades(context.Context, *QueryTradesRequest) (*QueryTradesResponse, error)
mustEmbedUnimplementedTradingServiceServer()
}
// UnimplementedTradingServiceServer must be embedded to have forward compatible implementations.
type UnimplementedTradingServiceServer struct {
}
func (UnimplementedTradingServiceServer) SubmitOrder(context.Context, *SubmitOrderRequest) (*SubmitOrderResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method SubmitOrder not implemented")
}
func (UnimplementedTradingServiceServer) CancelOrder(context.Context, *CancelOrderRequest) (*CancelOrderResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method CancelOrder not implemented")
}
func (UnimplementedTradingServiceServer) QueryOrder(context.Context, *QueryOrderRequest) (*QueryOrderResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method QueryOrder not implemented")
}
func (UnimplementedTradingServiceServer) QueryOrders(context.Context, *QueryOrdersRequest) (*QueryOrdersResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method QueryOrders not implemented")
}
func (UnimplementedTradingServiceServer) QueryTrades(context.Context, *QueryTradesRequest) (*QueryTradesResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method QueryTrades not implemented")
}
func (UnimplementedTradingServiceServer) mustEmbedUnimplementedTradingServiceServer() {}
// UnsafeTradingServiceServer may be embedded to opt out of forward compatibility for this service.
// Use of this interface is not recommended, as added methods to TradingServiceServer will
// result in compilation errors.
type UnsafeTradingServiceServer interface {
mustEmbedUnimplementedTradingServiceServer()
}
func RegisterTradingServiceServer(s grpc.ServiceRegistrar, srv TradingServiceServer) {
s.RegisterService(&TradingService_ServiceDesc, srv)
}
func _TradingService_SubmitOrder_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(SubmitOrderRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(TradingServiceServer).SubmitOrder(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: "/pb.TradingService/SubmitOrder",
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(TradingServiceServer).SubmitOrder(ctx, req.(*SubmitOrderRequest))
}
return interceptor(ctx, in, info, handler)
}
func _TradingService_CancelOrder_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(CancelOrderRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(TradingServiceServer).CancelOrder(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: "/pb.TradingService/CancelOrder",
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(TradingServiceServer).CancelOrder(ctx, req.(*CancelOrderRequest))
}
return interceptor(ctx, in, info, handler)
}
func _TradingService_QueryOrder_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(QueryOrderRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(TradingServiceServer).QueryOrder(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: "/pb.TradingService/QueryOrder",
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(TradingServiceServer).QueryOrder(ctx, req.(*QueryOrderRequest))
}
return interceptor(ctx, in, info, handler)
}
func _TradingService_QueryOrders_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(QueryOrdersRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(TradingServiceServer).QueryOrders(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: "/pb.TradingService/QueryOrders",
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(TradingServiceServer).QueryOrders(ctx, req.(*QueryOrdersRequest))
}
return interceptor(ctx, in, info, handler)
}
func _TradingService_QueryTrades_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(QueryTradesRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(TradingServiceServer).QueryTrades(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: "/pb.TradingService/QueryTrades",
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(TradingServiceServer).QueryTrades(ctx, req.(*QueryTradesRequest))
}
return interceptor(ctx, in, info, handler)
}
// TradingService_ServiceDesc is the grpc.ServiceDesc for TradingService service.
// It's only intended for direct use with grpc.RegisterService,
// and not to be introspected or modified (even as a copy)
var TradingService_ServiceDesc = grpc.ServiceDesc{
ServiceName: "pb.TradingService",
HandlerType: (*TradingServiceServer)(nil),
Methods: []grpc.MethodDesc{
{
MethodName: "SubmitOrder",
Handler: _TradingService_SubmitOrder_Handler,
},
{
MethodName: "CancelOrder",
Handler: _TradingService_CancelOrder_Handler,
},
{
MethodName: "QueryOrder",
Handler: _TradingService_QueryOrder_Handler,
},
{
MethodName: "QueryOrders",
Handler: _TradingService_QueryOrders_Handler,
},
{
MethodName: "QueryTrades",
Handler: _TradingService_QueryTrades_Handler,
},
},
Streams: []grpc.StreamDesc{},
Metadata: "pkg/pb/bbgo.proto",
}

File diff suppressed because one or more lines are too long

View File

@ -2,10 +2,10 @@
"""Client and server classes corresponding to protobuf-defined services.""" """Client and server classes corresponding to protobuf-defined services."""
import grpc import grpc
from . import bbgo_pb2 as bbgo__pb2 import bbgo_pb2 as bbgo__pb2
class BBGOStub(object): class MarketDataServiceStub(object):
"""Missing associated documentation comment in .proto file.""" """Missing associated documentation comment in .proto file."""
def __init__(self, channel): def __init__(self, channel):
@ -14,63 +14,191 @@ class BBGOStub(object):
Args: Args:
channel: A grpc.Channel. channel: A grpc.Channel.
""" """
self.Subcribe = channel.unary_stream( self.Subscribe = channel.unary_stream(
'/pb.BBGO/Subcribe', '/pb.MarketDataService/Subscribe',
request_serializer=bbgo__pb2.SubscribeRequest.SerializeToString, request_serializer=bbgo__pb2.SubscribeRequest.SerializeToString,
response_deserializer=bbgo__pb2.SubscribeResponse.FromString, response_deserializer=bbgo__pb2.SubscribeResponse.FromString,
) )
self.SubcribeUserData = channel.unary_stream(
'/pb.BBGO/SubcribeUserData',
request_serializer=bbgo__pb2.Empty.SerializeToString,
response_deserializer=bbgo__pb2.SubscribeResponse.FromString,
)
self.SubmitOrder = channel.unary_unary(
'/pb.BBGO/SubmitOrder',
request_serializer=bbgo__pb2.SubmitOrderRequest.SerializeToString,
response_deserializer=bbgo__pb2.SubmitOrderResponse.FromString,
)
self.CancelOrder = channel.unary_unary(
'/pb.BBGO/CancelOrder',
request_serializer=bbgo__pb2.CancelOrderRequest.SerializeToString,
response_deserializer=bbgo__pb2.CancelOrderResponse.FromString,
)
self.QueryOrder = channel.unary_unary(
'/pb.BBGO/QueryOrder',
request_serializer=bbgo__pb2.QueryOrderRequest.SerializeToString,
response_deserializer=bbgo__pb2.QueryOrderResponse.FromString,
)
self.QueryOrders = channel.unary_unary(
'/pb.BBGO/QueryOrders',
request_serializer=bbgo__pb2.QueryOrdersRequest.SerializeToString,
response_deserializer=bbgo__pb2.QueryOrdersResponse.FromString,
)
self.QueryTrades = channel.unary_unary(
'/pb.BBGO/QueryTrades',
request_serializer=bbgo__pb2.QueryTradesRequest.SerializeToString,
response_deserializer=bbgo__pb2.QueryTradesResponse.FromString,
)
self.QueryKLines = channel.unary_unary( self.QueryKLines = channel.unary_unary(
'/pb.BBGO/QueryKLines', '/pb.MarketDataService/QueryKLines',
request_serializer=bbgo__pb2.QueryKLinesRequest.SerializeToString, request_serializer=bbgo__pb2.QueryKLinesRequest.SerializeToString,
response_deserializer=bbgo__pb2.QueryKLinesResponse.FromString, response_deserializer=bbgo__pb2.QueryKLinesResponse.FromString,
) )
class BBGOServicer(object): class MarketDataServiceServicer(object):
"""Missing associated documentation comment in .proto file.""" """Missing associated documentation comment in .proto file."""
def Subcribe(self, request, context): def Subscribe(self, request, context):
"""Missing associated documentation comment in .proto file."""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def QueryKLines(self, request, context):
"""Missing associated documentation comment in .proto file."""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def add_MarketDataServiceServicer_to_server(servicer, server):
rpc_method_handlers = {
'Subscribe': grpc.unary_stream_rpc_method_handler(
servicer.Subscribe,
request_deserializer=bbgo__pb2.SubscribeRequest.FromString,
response_serializer=bbgo__pb2.SubscribeResponse.SerializeToString,
),
'QueryKLines': grpc.unary_unary_rpc_method_handler(
servicer.QueryKLines,
request_deserializer=bbgo__pb2.QueryKLinesRequest.FromString,
response_serializer=bbgo__pb2.QueryKLinesResponse.SerializeToString,
),
}
generic_handler = grpc.method_handlers_generic_handler(
'pb.MarketDataService', rpc_method_handlers)
server.add_generic_rpc_handlers((generic_handler,))
# This class is part of an EXPERIMENTAL API.
class MarketDataService(object):
"""Missing associated documentation comment in .proto file."""
@staticmethod
def Subscribe(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_stream(request, target, '/pb.MarketDataService/Subscribe',
bbgo__pb2.SubscribeRequest.SerializeToString,
bbgo__pb2.SubscribeResponse.FromString,
options, channel_credentials,
insecure, call_credentials, compression, wait_for_ready, timeout, metadata)
@staticmethod
def QueryKLines(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_unary(request, target, '/pb.MarketDataService/QueryKLines',
bbgo__pb2.QueryKLinesRequest.SerializeToString,
bbgo__pb2.QueryKLinesResponse.FromString,
options, channel_credentials,
insecure, call_credentials, compression, wait_for_ready, timeout, metadata)
class UserDataServiceStub(object):
"""Missing associated documentation comment in .proto file."""
def __init__(self, channel):
"""Constructor.
Args:
channel: A grpc.Channel.
"""
self.SubscribeUserData = channel.unary_stream(
'/pb.UserDataService/SubscribeUserData',
request_serializer=bbgo__pb2.Empty.SerializeToString,
response_deserializer=bbgo__pb2.SubscribeResponse.FromString,
)
class UserDataServiceServicer(object):
"""Missing associated documentation comment in .proto file."""
def SubscribeUserData(self, request, context):
"""should support streaming """should support streaming
""" """
context.set_code(grpc.StatusCode.UNIMPLEMENTED) context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!') context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!') raise NotImplementedError('Method not implemented!')
def SubcribeUserData(self, request, context):
def add_UserDataServiceServicer_to_server(servicer, server):
rpc_method_handlers = {
'SubscribeUserData': grpc.unary_stream_rpc_method_handler(
servicer.SubscribeUserData,
request_deserializer=bbgo__pb2.Empty.FromString,
response_serializer=bbgo__pb2.SubscribeResponse.SerializeToString,
),
}
generic_handler = grpc.method_handlers_generic_handler(
'pb.UserDataService', rpc_method_handlers)
server.add_generic_rpc_handlers((generic_handler,))
# This class is part of an EXPERIMENTAL API.
class UserDataService(object):
"""Missing associated documentation comment in .proto file."""
@staticmethod
def SubscribeUserData(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_stream(request, target, '/pb.UserDataService/SubscribeUserData',
bbgo__pb2.Empty.SerializeToString,
bbgo__pb2.SubscribeResponse.FromString,
options, channel_credentials,
insecure, call_credentials, compression, wait_for_ready, timeout, metadata)
class TradingServiceStub(object):
"""Missing associated documentation comment in .proto file."""
def __init__(self, channel):
"""Constructor.
Args:
channel: A grpc.Channel.
"""
self.SubmitOrder = channel.unary_unary(
'/pb.TradingService/SubmitOrder',
request_serializer=bbgo__pb2.SubmitOrderRequest.SerializeToString,
response_deserializer=bbgo__pb2.SubmitOrderResponse.FromString,
)
self.CancelOrder = channel.unary_unary(
'/pb.TradingService/CancelOrder',
request_serializer=bbgo__pb2.CancelOrderRequest.SerializeToString,
response_deserializer=bbgo__pb2.CancelOrderResponse.FromString,
)
self.QueryOrder = channel.unary_unary(
'/pb.TradingService/QueryOrder',
request_serializer=bbgo__pb2.QueryOrderRequest.SerializeToString,
response_deserializer=bbgo__pb2.QueryOrderResponse.FromString,
)
self.QueryOrders = channel.unary_unary(
'/pb.TradingService/QueryOrders',
request_serializer=bbgo__pb2.QueryOrdersRequest.SerializeToString,
response_deserializer=bbgo__pb2.QueryOrdersResponse.FromString,
)
self.QueryTrades = channel.unary_unary(
'/pb.TradingService/QueryTrades',
request_serializer=bbgo__pb2.QueryTradesRequest.SerializeToString,
response_deserializer=bbgo__pb2.QueryTradesResponse.FromString,
)
class TradingServiceServicer(object):
"""Missing associated documentation comment in .proto file.""" """Missing associated documentation comment in .proto file."""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def SubmitOrder(self, request, context): def SubmitOrder(self, request, context):
"""request-response """request-response
@ -103,25 +231,9 @@ class BBGOServicer(object):
context.set_details('Method not implemented!') context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!') raise NotImplementedError('Method not implemented!')
def QueryKLines(self, request, context):
"""Missing associated documentation comment in .proto file."""
context.set_code(grpc.StatusCode.UNIMPLEMENTED)
context.set_details('Method not implemented!')
raise NotImplementedError('Method not implemented!')
def add_TradingServiceServicer_to_server(servicer, server):
def add_BBGOServicer_to_server(servicer, server):
rpc_method_handlers = { rpc_method_handlers = {
'Subcribe': grpc.unary_stream_rpc_method_handler(
servicer.Subcribe,
request_deserializer=bbgo__pb2.SubscribeRequest.FromString,
response_serializer=bbgo__pb2.SubscribeResponse.SerializeToString,
),
'SubcribeUserData': grpc.unary_stream_rpc_method_handler(
servicer.SubcribeUserData,
request_deserializer=bbgo__pb2.Empty.FromString,
response_serializer=bbgo__pb2.SubscribeResponse.SerializeToString,
),
'SubmitOrder': grpc.unary_unary_rpc_method_handler( 'SubmitOrder': grpc.unary_unary_rpc_method_handler(
servicer.SubmitOrder, servicer.SubmitOrder,
request_deserializer=bbgo__pb2.SubmitOrderRequest.FromString, request_deserializer=bbgo__pb2.SubmitOrderRequest.FromString,
@ -147,55 +259,16 @@ def add_BBGOServicer_to_server(servicer, server):
request_deserializer=bbgo__pb2.QueryTradesRequest.FromString, request_deserializer=bbgo__pb2.QueryTradesRequest.FromString,
response_serializer=bbgo__pb2.QueryTradesResponse.SerializeToString, response_serializer=bbgo__pb2.QueryTradesResponse.SerializeToString,
), ),
'QueryKLines': grpc.unary_unary_rpc_method_handler(
servicer.QueryKLines,
request_deserializer=bbgo__pb2.QueryKLinesRequest.FromString,
response_serializer=bbgo__pb2.QueryKLinesResponse.SerializeToString,
),
} }
generic_handler = grpc.method_handlers_generic_handler( generic_handler = grpc.method_handlers_generic_handler(
'pb.BBGO', rpc_method_handlers) 'pb.TradingService', rpc_method_handlers)
server.add_generic_rpc_handlers((generic_handler,)) server.add_generic_rpc_handlers((generic_handler,))
# This class is part of an EXPERIMENTAL API. # This class is part of an EXPERIMENTAL API.
class BBGO(object): class TradingService(object):
"""Missing associated documentation comment in .proto file.""" """Missing associated documentation comment in .proto file."""
@staticmethod
def Subcribe(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_stream(request, target, '/pb.BBGO/Subcribe',
bbgo__pb2.SubscribeRequest.SerializeToString,
bbgo__pb2.SubscribeResponse.FromString,
options, channel_credentials,
insecure, call_credentials, compression, wait_for_ready, timeout, metadata)
@staticmethod
def SubcribeUserData(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_stream(request, target, '/pb.BBGO/SubcribeUserData',
bbgo__pb2.Empty.SerializeToString,
bbgo__pb2.SubscribeResponse.FromString,
options, channel_credentials,
insecure, call_credentials, compression, wait_for_ready, timeout, metadata)
@staticmethod @staticmethod
def SubmitOrder(request, def SubmitOrder(request,
target, target,
@ -207,7 +280,7 @@ class BBGO(object):
wait_for_ready=None, wait_for_ready=None,
timeout=None, timeout=None,
metadata=None): metadata=None):
return grpc.experimental.unary_unary(request, target, '/pb.BBGO/SubmitOrder', return grpc.experimental.unary_unary(request, target, '/pb.TradingService/SubmitOrder',
bbgo__pb2.SubmitOrderRequest.SerializeToString, bbgo__pb2.SubmitOrderRequest.SerializeToString,
bbgo__pb2.SubmitOrderResponse.FromString, bbgo__pb2.SubmitOrderResponse.FromString,
options, channel_credentials, options, channel_credentials,
@ -224,7 +297,7 @@ class BBGO(object):
wait_for_ready=None, wait_for_ready=None,
timeout=None, timeout=None,
metadata=None): metadata=None):
return grpc.experimental.unary_unary(request, target, '/pb.BBGO/CancelOrder', return grpc.experimental.unary_unary(request, target, '/pb.TradingService/CancelOrder',
bbgo__pb2.CancelOrderRequest.SerializeToString, bbgo__pb2.CancelOrderRequest.SerializeToString,
bbgo__pb2.CancelOrderResponse.FromString, bbgo__pb2.CancelOrderResponse.FromString,
options, channel_credentials, options, channel_credentials,
@ -241,7 +314,7 @@ class BBGO(object):
wait_for_ready=None, wait_for_ready=None,
timeout=None, timeout=None,
metadata=None): metadata=None):
return grpc.experimental.unary_unary(request, target, '/pb.BBGO/QueryOrder', return grpc.experimental.unary_unary(request, target, '/pb.TradingService/QueryOrder',
bbgo__pb2.QueryOrderRequest.SerializeToString, bbgo__pb2.QueryOrderRequest.SerializeToString,
bbgo__pb2.QueryOrderResponse.FromString, bbgo__pb2.QueryOrderResponse.FromString,
options, channel_credentials, options, channel_credentials,
@ -258,7 +331,7 @@ class BBGO(object):
wait_for_ready=None, wait_for_ready=None,
timeout=None, timeout=None,
metadata=None): metadata=None):
return grpc.experimental.unary_unary(request, target, '/pb.BBGO/QueryOrders', return grpc.experimental.unary_unary(request, target, '/pb.TradingService/QueryOrders',
bbgo__pb2.QueryOrdersRequest.SerializeToString, bbgo__pb2.QueryOrdersRequest.SerializeToString,
bbgo__pb2.QueryOrdersResponse.FromString, bbgo__pb2.QueryOrdersResponse.FromString,
options, channel_credentials, options, channel_credentials,
@ -275,25 +348,8 @@ class BBGO(object):
wait_for_ready=None, wait_for_ready=None,
timeout=None, timeout=None,
metadata=None): metadata=None):
return grpc.experimental.unary_unary(request, target, '/pb.BBGO/QueryTrades', return grpc.experimental.unary_unary(request, target, '/pb.TradingService/QueryTrades',
bbgo__pb2.QueryTradesRequest.SerializeToString, bbgo__pb2.QueryTradesRequest.SerializeToString,
bbgo__pb2.QueryTradesResponse.FromString, bbgo__pb2.QueryTradesResponse.FromString,
options, channel_credentials, options, channel_credentials,
insecure, call_credentials, compression, wait_for_ready, timeout, metadata) insecure, call_credentials, compression, wait_for_ready, timeout, metadata)
@staticmethod
def QueryKLines(request,
target,
options=(),
channel_credentials=None,
call_credentials=None,
insecure=False,
compression=None,
wait_for_ready=None,
timeout=None,
metadata=None):
return grpc.experimental.unary_unary(request, target, '/pb.BBGO/QueryKLines',
bbgo__pb2.QueryKLinesRequest.SerializeToString,
bbgo__pb2.QueryKLinesResponse.FromString,
options, channel_credentials,
insecure, call_credentials, compression, wait_for_ready, timeout, metadata)