feat(server/v2): add config overwrite (#20874)
This commit is contained in:
@@ -0,0 +1,95 @@
|
||||
package grpc
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
|
||||
gogoproto "github.com/cosmos/gogoproto/proto"
|
||||
"google.golang.org/grpc/encoding"
|
||||
"google.golang.org/protobuf/proto"
|
||||
|
||||
_ "cosmossdk.io/api/amino" // Import amino.proto file for reflection
|
||||
appmanager "cosmossdk.io/core/app"
|
||||
)
|
||||
|
||||
type protoCodec struct {
|
||||
interfaceRegistry appmanager.InterfaceRegistry
|
||||
}
|
||||
|
||||
// newProtoCodec returns a reference to a new ProtoCodec
|
||||
func newProtoCodec(interfaceRegistry appmanager.InterfaceRegistry) *protoCodec {
|
||||
return &protoCodec{
|
||||
interfaceRegistry: interfaceRegistry,
|
||||
}
|
||||
}
|
||||
|
||||
// Marshal implements BinaryMarshaler.Marshal method.
|
||||
// NOTE: this function must be used with a concrete type which
|
||||
// implements proto.Message. For interface please use the codec.MarshalInterface
|
||||
func (pc *protoCodec) Marshal(o gogoproto.Message) ([]byte, error) {
|
||||
// Size() check can catch the typed nil value.
|
||||
if o == nil || gogoproto.Size(o) == 0 {
|
||||
// return empty bytes instead of nil, because nil has special meaning in places like store.Set
|
||||
return []byte{}, nil
|
||||
}
|
||||
|
||||
return gogoproto.Marshal(o)
|
||||
}
|
||||
|
||||
// Unmarshal implements BinaryMarshaler.Unmarshal method.
|
||||
// NOTE: this function must be used with a concrete type which
|
||||
// implements proto.Message. For interface please use the codec.UnmarshalInterface
|
||||
func (pc *protoCodec) Unmarshal(bz []byte, ptr gogoproto.Message) error {
|
||||
err := gogoproto.Unmarshal(bz, ptr)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// err = codectypes.UnpackInterfaces(ptr, pc.interfaceRegistry) // TODO: identify if needed for grpc
|
||||
// if err != nil {
|
||||
// return err
|
||||
// }
|
||||
return nil
|
||||
}
|
||||
|
||||
func (pc *protoCodec) Name() string {
|
||||
return "cosmos-sdk-grpc-codec"
|
||||
}
|
||||
|
||||
// GRPCCodec returns the gRPC Codec for this specific ProtoCodec
|
||||
func (pc *protoCodec) GRPCCodec() encoding.Codec {
|
||||
return &grpcProtoCodec{cdc: pc}
|
||||
}
|
||||
|
||||
// grpcProtoCodec is the implementation of the gRPC proto codec.
|
||||
type grpcProtoCodec struct {
|
||||
cdc appmanager.ProtoCodec
|
||||
}
|
||||
|
||||
var errUnknownProtoType = errors.New("codec: unknown proto type") // sentinel error
|
||||
|
||||
func (c grpcProtoCodec) Marshal(v any) ([]byte, error) {
|
||||
switch m := v.(type) {
|
||||
case proto.Message:
|
||||
protov2MarshalOpts := proto.MarshalOptions{Deterministic: true}
|
||||
return protov2MarshalOpts.Marshal(m)
|
||||
case gogoproto.Message:
|
||||
return c.cdc.Marshal(m)
|
||||
default:
|
||||
return nil, fmt.Errorf("%w: cannot marshal type %T", errUnknownProtoType, v)
|
||||
}
|
||||
}
|
||||
|
||||
func (c grpcProtoCodec) Unmarshal(data []byte, v any) error {
|
||||
switch m := v.(type) {
|
||||
case proto.Message:
|
||||
return proto.Unmarshal(data, m)
|
||||
case gogoproto.Message:
|
||||
return c.cdc.Unmarshal(data, m)
|
||||
default:
|
||||
return fmt.Errorf("%w: cannot unmarshal type %T", errUnknownProtoType, v)
|
||||
}
|
||||
}
|
||||
|
||||
func (c grpcProtoCodec) Name() string {
|
||||
return "cosmos-sdk-grpc-codec"
|
||||
}
|
||||
@@ -32,3 +32,20 @@ type Config struct {
|
||||
// The default value is math.MaxInt32.
|
||||
MaxSendMsgSize int `mapstructure:"max-send-msg-size" toml:"max-send-msg-size" comment:"MaxSendMsgSize defines the max message size in bytes the server can send.\nThe default value is math.MaxInt32."`
|
||||
}
|
||||
|
||||
// CfgOption is a function that allows to overwrite the default server configuration.
|
||||
type CfgOption func(*Config)
|
||||
|
||||
// OverwriteDefaultConfig overwrites the default config with the new config.
|
||||
func OverwriteDefaultConfig(newCfg *Config) CfgOption {
|
||||
return func(cfg *Config) {
|
||||
*cfg = *newCfg
|
||||
}
|
||||
}
|
||||
|
||||
// Disable the grpc-gateway server by default (default enabled).
|
||||
func Disable() CfgOption {
|
||||
return func(cfg *Config) {
|
||||
cfg.Enable = false
|
||||
}
|
||||
}
|
||||
|
||||
+47
-124
@@ -2,19 +2,12 @@ package grpc
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net"
|
||||
|
||||
gogogrpc "github.com/cosmos/gogoproto/grpc"
|
||||
gogoproto "github.com/cosmos/gogoproto/proto"
|
||||
"github.com/spf13/viper"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/encoding"
|
||||
"google.golang.org/protobuf/proto"
|
||||
|
||||
_ "cosmossdk.io/api/amino" // Import amino.proto file for reflection
|
||||
appmanager "cosmossdk.io/core/app"
|
||||
"cosmossdk.io/core/transaction"
|
||||
"cosmossdk.io/log"
|
||||
serverv2 "cosmossdk.io/server/v2"
|
||||
@@ -22,27 +15,26 @@ import (
|
||||
)
|
||||
|
||||
type GRPCServer[AppT serverv2.AppI[T], T transaction.Tx] struct {
|
||||
logger log.Logger
|
||||
config *Config
|
||||
logger log.Logger
|
||||
config *Config
|
||||
cfgOptions []CfgOption
|
||||
|
||||
grpcSrv *grpc.Server
|
||||
}
|
||||
|
||||
type GRPCService interface {
|
||||
// RegisterGRPCServer registers gRPC services directly with the gRPC server.
|
||||
RegisterGRPCServer(gogogrpc.Server)
|
||||
}
|
||||
|
||||
func New[AppT serverv2.AppI[T], T transaction.Tx]() *GRPCServer[AppT, T] {
|
||||
return &GRPCServer[AppT, T]{}
|
||||
// New creates a new grpc server.
|
||||
func New[AppT serverv2.AppI[T], T transaction.Tx](cfgOptions ...CfgOption) *GRPCServer[AppT, T] {
|
||||
return &GRPCServer[AppT, T]{
|
||||
cfgOptions: cfgOptions,
|
||||
}
|
||||
}
|
||||
|
||||
// Init returns a correctly configured and initialized gRPC server.
|
||||
// Note, the caller is responsible for starting the server.
|
||||
func (g *GRPCServer[AppT, T]) Init(appI AppT, v *viper.Viper, logger log.Logger) error {
|
||||
cfg := DefaultConfig()
|
||||
func (s *GRPCServer[AppT, T]) Init(appI AppT, v *viper.Viper, logger log.Logger) error {
|
||||
cfg := s.Config().(*Config)
|
||||
if v != nil {
|
||||
if err := v.Sub(g.Name()).Unmarshal(&cfg); err != nil {
|
||||
if err := v.Sub(s.Name()).Unmarshal(&cfg); err != nil {
|
||||
return fmt.Errorf("failed to unmarshal config: %w", err)
|
||||
}
|
||||
}
|
||||
@@ -55,25 +47,42 @@ func (g *GRPCServer[AppT, T]) Init(appI AppT, v *viper.Viper, logger log.Logger)
|
||||
|
||||
// appI.RegisterGRPCServer(grpcSrv)
|
||||
|
||||
// Reflection allows external clients to see what services and methods
|
||||
// the gRPC server exposes.
|
||||
// Reflection allows external clients to see what services and methods the gRPC server exposes.
|
||||
gogoreflection.Register(grpcSrv)
|
||||
|
||||
g.grpcSrv = grpcSrv
|
||||
g.config = cfg
|
||||
g.logger = logger.With(log.ModuleKey, g.Name())
|
||||
s.grpcSrv = grpcSrv
|
||||
s.config = cfg
|
||||
s.logger = logger.With(log.ModuleKey, s.Name())
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer[AppT, T]) Name() string {
|
||||
func (s *GRPCServer[AppT, T]) Name() string {
|
||||
return "grpc"
|
||||
}
|
||||
|
||||
func (g *GRPCServer[AppT, T]) Start(ctx context.Context) error {
|
||||
listener, err := net.Listen("tcp", g.config.Address)
|
||||
func (s *GRPCServer[AppT, T]) Config() any {
|
||||
if s.config == nil || s.config == (&Config{}) {
|
||||
cfg := DefaultConfig()
|
||||
// overwrite the default config with the provided options
|
||||
for _, opt := range s.cfgOptions {
|
||||
opt(cfg)
|
||||
}
|
||||
|
||||
return cfg
|
||||
}
|
||||
|
||||
return s.config
|
||||
}
|
||||
|
||||
func (s *GRPCServer[AppT, T]) Start(ctx context.Context) error {
|
||||
if !s.config.Enable {
|
||||
return nil
|
||||
}
|
||||
|
||||
listener, err := net.Listen("tcp", s.config.Address)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to listen on address %s: %w", g.config.Address, err)
|
||||
return fmt.Errorf("failed to listen on address %s: %w", s.config.Address, err)
|
||||
}
|
||||
|
||||
errCh := make(chan error)
|
||||
@@ -82,110 +91,24 @@ func (g *GRPCServer[AppT, T]) Start(ctx context.Context) error {
|
||||
// an error upon failure, which we'll send on the error channel that will be
|
||||
// consumed by the for block below.
|
||||
go func() {
|
||||
g.logger.Info("starting gRPC server...", "address", g.config.Address)
|
||||
errCh <- g.grpcSrv.Serve(listener)
|
||||
s.logger.Info("starting gRPC server...", "address", s.config.Address)
|
||||
errCh <- s.grpcSrv.Serve(listener)
|
||||
}()
|
||||
|
||||
// Start a blocking select to wait for an indication to stop the server or that
|
||||
// the server failed to start properly.
|
||||
err = <-errCh
|
||||
g.logger.Error("failed to start gRPC server", "err", err)
|
||||
s.logger.Error("failed to start gRPC server", "err", err)
|
||||
return err
|
||||
}
|
||||
|
||||
func (g *GRPCServer[AppT, T]) Stop(ctx context.Context) error {
|
||||
g.logger.Info("stopping gRPC server...", "address", g.config.Address)
|
||||
g.grpcSrv.GracefulStop()
|
||||
func (s *GRPCServer[AppT, T]) Stop(ctx context.Context) error {
|
||||
if !s.config.Enable {
|
||||
return nil
|
||||
}
|
||||
|
||||
s.logger.Info("stopping gRPC server...", "address", s.config.Address)
|
||||
s.grpcSrv.GracefulStop()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (g *GRPCServer[AppT, T]) Config() any {
|
||||
if g.config == nil || g.config == (&Config{}) {
|
||||
return DefaultConfig()
|
||||
}
|
||||
|
||||
return g.config
|
||||
}
|
||||
|
||||
type protoCodec struct {
|
||||
interfaceRegistry appmanager.InterfaceRegistry
|
||||
}
|
||||
|
||||
// newProtoCodec returns a reference to a new ProtoCodec
|
||||
func newProtoCodec(interfaceRegistry appmanager.InterfaceRegistry) *protoCodec {
|
||||
return &protoCodec{
|
||||
interfaceRegistry: interfaceRegistry,
|
||||
}
|
||||
}
|
||||
|
||||
// Marshal implements BinaryMarshaler.Marshal method.
|
||||
// NOTE: this function must be used with a concrete type which
|
||||
// implements proto.Message. For interface please use the codec.MarshalInterface
|
||||
func (pc *protoCodec) Marshal(o gogoproto.Message) ([]byte, error) {
|
||||
// Size() check can catch the typed nil value.
|
||||
if o == nil || gogoproto.Size(o) == 0 {
|
||||
// return empty bytes instead of nil, because nil has special meaning in places like store.Set
|
||||
return []byte{}, nil
|
||||
}
|
||||
|
||||
return gogoproto.Marshal(o)
|
||||
}
|
||||
|
||||
// Unmarshal implements BinaryMarshaler.Unmarshal method.
|
||||
// NOTE: this function must be used with a concrete type which
|
||||
// implements proto.Message. For interface please use the codec.UnmarshalInterface
|
||||
func (pc *protoCodec) Unmarshal(bz []byte, ptr gogoproto.Message) error {
|
||||
err := gogoproto.Unmarshal(bz, ptr)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// err = codectypes.UnpackInterfaces(ptr, pc.interfaceRegistry) // TODO: identify if needed for grpc
|
||||
// if err != nil {
|
||||
// return err
|
||||
// }
|
||||
return nil
|
||||
}
|
||||
|
||||
func (pc *protoCodec) Name() string {
|
||||
return "cosmos-sdk-grpc-codec"
|
||||
}
|
||||
|
||||
// GRPCCodec returns the gRPC Codec for this specific ProtoCodec
|
||||
func (pc *protoCodec) GRPCCodec() encoding.Codec {
|
||||
return &grpcProtoCodec{cdc: pc}
|
||||
}
|
||||
|
||||
// grpcProtoCodec is the implementation of the gRPC proto codec.
|
||||
type grpcProtoCodec struct {
|
||||
cdc appmanager.ProtoCodec
|
||||
}
|
||||
|
||||
var errUnknownProtoType = errors.New("codec: unknown proto type") // sentinel error
|
||||
|
||||
func (g grpcProtoCodec) Marshal(v any) ([]byte, error) {
|
||||
switch m := v.(type) {
|
||||
case proto.Message:
|
||||
protov2MarshalOpts := proto.MarshalOptions{Deterministic: true}
|
||||
return protov2MarshalOpts.Marshal(m)
|
||||
case gogoproto.Message:
|
||||
return g.cdc.Marshal(m)
|
||||
default:
|
||||
return nil, fmt.Errorf("%w: cannot marshal type %T", errUnknownProtoType, v)
|
||||
}
|
||||
}
|
||||
|
||||
func (g grpcProtoCodec) Unmarshal(data []byte, v any) error {
|
||||
switch m := v.(type) {
|
||||
case proto.Message:
|
||||
return proto.Unmarshal(data, m)
|
||||
case gogoproto.Message:
|
||||
return g.cdc.Unmarshal(data, m)
|
||||
default:
|
||||
return fmt.Errorf("%w: cannot unmarshal type %T", errUnknownProtoType, v)
|
||||
}
|
||||
}
|
||||
|
||||
func (g grpcProtoCodec) Name() string {
|
||||
return "cosmos-sdk-grpc-codec"
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user