+2
-20
@@ -297,9 +297,6 @@ func (w *Watcher) TrapSignals() {
|
||||
}
|
||||
}
|
||||
|
||||
// TODO: Ideally this would call exit() but properly return an error. The
|
||||
// exit() is problematic (i.e. racey) especiaily when orchestrating multiple
|
||||
// reva services from some external runtime (like in the "opencloud server" case
|
||||
func gracefulShutdown(w *Watcher) {
|
||||
defer w.Clean()
|
||||
w.log.Info().Int("Timeout", w.gracefulShutdownTimeout).Msg("preparing for a graceful shutdown with deadline")
|
||||
@@ -309,7 +306,7 @@ func gracefulShutdown(w *Watcher) {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
w.log.Info().Msgf("fd to %s:%s gracefully closed", s.Network(), s.Address())
|
||||
w.log.Info().Str("network.transport", s.Network()).Str("network.local.address", s.Address()).Msg("fd gracefully closed")
|
||||
err := s.GracefulStop()
|
||||
if err != nil {
|
||||
w.log.Error().Err(err).Msg("error stopping server")
|
||||
@@ -327,7 +324,7 @@ func gracefulShutdown(w *Watcher) {
|
||||
case <-time.After(time.Duration(w.gracefulShutdownTimeout) * time.Second):
|
||||
w.log.Info().Msg("graceful shutdown timeout reached. running hard shutdown")
|
||||
for _, s := range w.ss {
|
||||
w.log.Info().Msgf("fd to %s:%s abruptly closed", s.Network(), s.Address())
|
||||
w.log.Info().Str("network.transport", s.Network()).Str("network.local.address", s.Address()).Msg("fd abruptly closed")
|
||||
err := s.Stop()
|
||||
if err != nil {
|
||||
w.log.Error().Err(err).Msg("error stopping server")
|
||||
@@ -340,21 +337,6 @@ func gracefulShutdown(w *Watcher) {
|
||||
}
|
||||
}
|
||||
|
||||
// TODO: Ideally this would call exit() but properly return an error. The
|
||||
// exit() is problematic (i.e. racey) especiaily when orchestrating multiple
|
||||
// reva services from some external runtime (like in the "opencloud server" case
|
||||
func hardShutdown(w *Watcher) {
|
||||
w.log.Info().Msg("preparing for hard shutdown, aborting all conns")
|
||||
for _, s := range w.ss {
|
||||
w.log.Info().Msgf("fd to %s:%s abruptly closed", s.Network(), s.Address())
|
||||
err := s.Stop()
|
||||
if err != nil {
|
||||
w.log.Error().Err(err).Msg("error stopping server")
|
||||
}
|
||||
}
|
||||
w.Exit(0)
|
||||
}
|
||||
|
||||
func getListenerFile(ln net.Listener) (*os.File, error) {
|
||||
switch t := ln.(type) {
|
||||
case *net.TCPListener:
|
||||
|
||||
+50
-42
@@ -45,16 +45,18 @@ type server struct {
|
||||
// Returns nil if no http server is configured in the config file.
|
||||
// The GracefulShutdownTimeout set to default 20 seconds and can be overridden in the core config.
|
||||
// Logging a fatal error and exit with code 1 if the http server cannot be created.
|
||||
func NewDrivenHTTPServerWithOptions(mainConf map[string]interface{}, opts ...Option) RevaDrivenServer {
|
||||
func NewDrivenHTTPServerWithOptions(mainConf map[string]interface{}, opts ...Option) *server {
|
||||
if !isEnabledHTTP(mainConf) {
|
||||
return nil
|
||||
}
|
||||
options := newOptions(opts...)
|
||||
if srv := newServer(HTTP, mainConf, options); srv != nil {
|
||||
return srv
|
||||
var srv *server
|
||||
var err error
|
||||
if srv, err = newServer(HTTP, mainConf, options); err != nil {
|
||||
options.Logger.Fatal().Err(err).Msg("failed to create http server")
|
||||
}
|
||||
options.Logger.Fatal().Msg("nothing to do, no http enabled_services declared in config")
|
||||
return nil
|
||||
return srv
|
||||
|
||||
}
|
||||
|
||||
// NewDrivenGRPCServerWithOptions runs a revad server w/o watcher with the given config file and options.
|
||||
@@ -62,36 +64,37 @@ func NewDrivenHTTPServerWithOptions(mainConf map[string]interface{}, opts ...Opt
|
||||
// Returns nil if no grpc server is configured in the config file.
|
||||
// The GracefulShutdownTimeout set to default 20 seconds and can be overridden in the core config.
|
||||
// Logging a fatal error and exit with code 1 if the grpc server cannot be created.
|
||||
func NewDrivenGRPCServerWithOptions(mainConf map[string]interface{}, opts ...Option) RevaDrivenServer {
|
||||
func NewDrivenGRPCServerWithOptions(mainConf map[string]interface{}, opts ...Option) *server {
|
||||
if !isEnabledGRPC(mainConf) {
|
||||
return nil
|
||||
}
|
||||
options := newOptions(opts...)
|
||||
if srv := newServer(GRPC, mainConf, options); srv != nil {
|
||||
return srv
|
||||
var srv *server
|
||||
var err error
|
||||
if srv, err = newServer(GRPC, mainConf, options); err != nil {
|
||||
options.Logger.Fatal().Err(err).Msg("failed to create grpc server")
|
||||
}
|
||||
options.Logger.Fatal().Msg("nothing to do, no grpc enabled_services declared in config")
|
||||
return nil
|
||||
return srv
|
||||
}
|
||||
|
||||
// Start starts the reva server, listening on the configured address and network.
|
||||
func (s *server) Start() error {
|
||||
if s.srv == nil {
|
||||
err := fmt.Errorf("reva %s server not initialized", s.protocol)
|
||||
s.log.Fatal().Err(err).Send()
|
||||
return err
|
||||
return fmt.Errorf("reva %s server not initialized", s.protocol)
|
||||
}
|
||||
ln, err := net.Listen(s.srv.Network(), s.srv.Address())
|
||||
if err != nil {
|
||||
s.log.Fatal().Err(err).Send()
|
||||
return err
|
||||
}
|
||||
if err = s.srv.Start(ln); err != nil {
|
||||
if !errors.Is(err, http.ErrServerClosed) {
|
||||
s.log.Error().Err(err).Msgf("reva %s server error", s.protocol)
|
||||
s.log.Error().Err(err).Msg("reva server error")
|
||||
}
|
||||
return err
|
||||
}
|
||||
// update logger with transport and address
|
||||
logger := s.log.With().Str("network.transport", s.srv.Network()).Str("network.local.address", s.srv.Address()).Logger()
|
||||
s.log = &logger
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -102,10 +105,13 @@ func (s *server) Stop() error {
|
||||
}
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
s.log.Info().Msgf("gracefully stopping %s:%s reva %s server", s.srv.Network(), s.srv.Address(), s.protocol)
|
||||
s.log.Info().Msg("gracefully stopping reva server")
|
||||
if err := s.srv.GracefulStop(); err != nil {
|
||||
s.log.Error().Err(err).Msgf("error gracefully stopping reva %s server", s.protocol)
|
||||
s.srv.Stop()
|
||||
s.log.Error().Err(err).Msg("error gracefully stopping reva server")
|
||||
err := s.srv.Stop()
|
||||
if err != nil {
|
||||
s.log.Error().Err(err).Msg("error stopping reva server")
|
||||
}
|
||||
}
|
||||
close(done)
|
||||
}()
|
||||
@@ -115,64 +121,66 @@ func (s *server) Stop() error {
|
||||
s.log.Info().Msg("graceful shutdown timeout reached. running hard shutdown")
|
||||
err := s.srv.Stop()
|
||||
if err != nil {
|
||||
s.log.Error().Err(err).Msgf("error stopping reva %s server", s.protocol)
|
||||
s.log.Error().Err(err).Msg("error stopping reva server")
|
||||
}
|
||||
return nil
|
||||
case <-done:
|
||||
s.log.Info().Msgf("reva %s server gracefully stopped", s.protocol)
|
||||
s.log.Info().Msg("reva server gracefully stopped")
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// newServer runs a revad server w/o watcher with the given config file and options.
|
||||
func newServer(protocol int, mainConf map[string]interface{}, options Options) RevaDrivenServer {
|
||||
func newServer(protocol int, mainConf map[string]interface{}, options Options) (*server, error) {
|
||||
parseSharedConfOrDie(mainConf["shared"])
|
||||
coreConf := parseCoreConfOrDie(mainConf["core"])
|
||||
log := options.Logger
|
||||
|
||||
if err := registry.Init(options.Registry); err != nil {
|
||||
log.Fatal().Err(err).Msg("failed to initialize registry client")
|
||||
return nil
|
||||
return nil, err
|
||||
}
|
||||
|
||||
srv := &server{}
|
||||
|
||||
// update logger with hostname
|
||||
host, _ := os.Hostname()
|
||||
log.Info().Msgf("host info: %s", host)
|
||||
logger := options.Logger.With().Str("host.name", host).Logger()
|
||||
srv.log = &logger
|
||||
|
||||
// Only initialize tracing if we didn't get a tracer provider.
|
||||
if options.TraceProvider == nil {
|
||||
log.Debug().Msg("no pre-existing tracer given, initializing tracing")
|
||||
srv.log.Debug().Msg("no pre-existing tracer given, initializing tracing")
|
||||
options.TraceProvider = initTracing(coreConf)
|
||||
}
|
||||
initCPUCount(coreConf, log)
|
||||
initCPUCount(coreConf, srv.log)
|
||||
|
||||
gracefulShutdownTimeout := 20 * time.Second
|
||||
srv.gracefulShutdownTimeout = 20 * time.Second
|
||||
if coreConf.GracefulShutdownTimeout > 0 {
|
||||
gracefulShutdownTimeout = time.Duration(coreConf.GracefulShutdownTimeout) * time.Second
|
||||
srv.gracefulShutdownTimeout = time.Duration(coreConf.GracefulShutdownTimeout) * time.Second
|
||||
}
|
||||
|
||||
srv := &server{
|
||||
log: options.Logger,
|
||||
gracefulShutdownTimeout: gracefulShutdownTimeout,
|
||||
}
|
||||
switch protocol {
|
||||
case HTTP:
|
||||
s, err := getHTTPServer(mainConf["http"], options.Logger, options.TraceProvider)
|
||||
s, err := getHTTPServer(mainConf["http"], srv.log, options.TraceProvider)
|
||||
if err != nil {
|
||||
options.Logger.Fatal().Err(err).Msg("error creating http server")
|
||||
return nil
|
||||
return nil, err
|
||||
}
|
||||
srv.srv = s
|
||||
srv.protocol = "http"
|
||||
return srv
|
||||
// update logger with protocol
|
||||
logger := srv.log.With().Str("protocol", "http").Logger()
|
||||
srv.log = &logger
|
||||
return srv, nil
|
||||
case GRPC:
|
||||
s, err := getGRPCServer(mainConf["grpc"], options.Logger, options.TraceProvider)
|
||||
s, err := getGRPCServer(mainConf["grpc"], srv.log, options.TraceProvider)
|
||||
if err != nil {
|
||||
options.Logger.Fatal().Err(err).Msg("error creating grpc server")
|
||||
return nil
|
||||
return nil, err
|
||||
}
|
||||
srv.srv = s
|
||||
srv.protocol = "grpc"
|
||||
return srv
|
||||
// update logger with protocol
|
||||
logger := srv.log.With().Str("protocol", "grpc").Logger()
|
||||
srv.log = &logger
|
||||
return srv, nil
|
||||
}
|
||||
return nil
|
||||
return nil, fmt.Errorf("unknown protocol: %d", protocol)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user