/
glukhov2008
/
device-api
Обзор
Документация
Войти
/
glukhov2008
/
device-api
Код
Запросы
0
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
internal/server/server.go
172 строки
5 KB
AndrewGluss
other folders
25 мар 2024, 15:44
25 мар 2024, 15:44
2bc35c2
Код
Авторство
О чём код?
package server import ( "context" "errors" "fmt" repo2 "gitlab.ozon.dev/qa/classroom-4/act-device-api/internal/app/repo" "net" "net/http" "os" "os/signal" "sync/atomic" "syscall" "time" "github.com/jmoiron/sqlx" "github.com/rs/zerolog/log" "google.golang.org/grpc" "google.golang.org/grpc/keepalive" "google.golang.org/grpc/reflection" grpc_middleware "github.com/grpc-ecosystem/go-grpc-middleware" grpcrecovery "github.com/grpc-ecosystem/go-grpc-middleware/recovery" grpc_ctxtags "github.com/grpc-ecosystem/go-grpc-middleware/tags" grpc_opentracing "github.com/grpc-ecosystem/go-grpc-middleware/tracing/opentracing" grpc_prometheus "github.com/grpc-ecosystem/go-grpc-prometheus" "gitlab.ozon.dev/qa/classroom-4/act-device-api/internal/api" "gitlab.ozon.dev/qa/classroom-4/act-device-api/internal/config" pb "gitlab.ozon.dev/qa/classroom-4/act-device-api/pkg/act-device-api/gitlab.ozon.dev/qa/classroom-4/act-device-api/pkg/act-device-api" ) // GrpcServer is gRPC server type GrpcServer struct { db *sqlx.DB batchSize uint } // NewGrpcServer returns gRPC server with supporting of batch listing func NewGrpcServer(db *sqlx.DB, batchSize uint) *GrpcServer { return &GrpcServer{ db: db, batchSize: batchSize, } } // Start method runs server func (s *GrpcServer) Start(cfg *config.Config) error { ctx, cancel := context.WithCancel(context.Background()) defer cancel() gatewayAddr := fmt.Sprintf("%s:%v", cfg.Rest.Host, cfg.Rest.Port) grpcAddr := fmt.Sprintf("%s:%v", cfg.Grpc.Host, cfg.Grpc.Port) metricsAddr := fmt.Sprintf("%s:%v", cfg.Metrics.Host, cfg.Metrics.Port) gatewayServer := createGatewayServer(grpcAddr, gatewayAddr) go func() { log.Info().Msgf("Gateway server is running on %s", gatewayAddr) if err := gatewayServer.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { log.Error().Err(err).Msg("Failed running gateway server") cancel() } }() metricsServer := createMetricsServer(cfg) go func() { log.Info().Msgf("Metrics server is running on %s", metricsAddr) if err := metricsServer.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { log.Error().Err(err).Msg("Failed running metrics server") cancel() } }() isReady := &atomic.Value{} isReady.Store(false) statusServer := createStatusServer(cfg, isReady) go func() { statusAdrr := fmt.Sprintf("%s:%v", cfg.Status.Host, cfg.Status.Port) log.Info().Msgf("Status server is running on %s", statusAdrr) if err := statusServer.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { log.Error().Err(err).Msg("Failed running status server") } }() l, err := net.Listen("tcp", grpcAddr) if err != nil { return fmt.Errorf("failed to listen: %w", err) } defer l.Close() grpcServer := grpc.NewServer( grpc.KeepaliveParams(keepalive.ServerParameters{ MaxConnectionIdle: time.Duration(cfg.Grpc.MaxConnectionIdle) * time.Minute, Timeout: time.Duration(cfg.Grpc.Timeout) * time.Second, MaxConnectionAge: time.Duration(cfg.Grpc.MaxConnectionAge) * time.Minute, Time: time.Duration(cfg.Grpc.Timeout) * time.Minute, }), grpc.UnaryInterceptor(grpc_middleware.ChainUnaryServer( grpc_ctxtags.UnaryServerInterceptor(), grpc_prometheus.UnaryServerInterceptor, grpc_opentracing.UnaryServerInterceptor(), grpcrecovery.UnaryServerInterceptor(), RequestLogInterceptor(), ResponseLogInterceptor(), )), ) r := repo2.NewRepo(s.db, s.batchSize) er := repo2.NewEventRepo(s.db, s.batchSize) notificationRepo := repo2.NewNotificationRepo(s.db, s.batchSize) pb.RegisterActDeviceApiServiceServer(grpcServer, api.NewDeviceAPI(r, er)) pb.RegisterActNotificationApiServiceServer(grpcServer, api.NewNotificationAPI(notificationRepo)) grpc_prometheus.EnableHandlingTimeHistogram() grpc_prometheus.Register(grpcServer) go func() { log.Info().Msgf("GRPC Server is listening on: %s", grpcAddr) if err := grpcServer.Serve(l); err != nil { log.Fatal().Err(err).Msg("Failed running gRPC server") } }() go func() { time.Sleep(2 * time.Second) isReady.Store(true) log.Info().Msg("The service is ready to accept requests") }() if cfg.Project.Debug { reflection.Register(grpcServer) } quit := make(chan os.Signal, 1) signal.Notify(quit, os.Interrupt, syscall.SIGTERM) select { case v := <-quit: log.Info().Msgf("signal.Notify: %v", v) case done := <-ctx.Done(): log.Info().Msgf("ctx.Done: %v", done) } isReady.Store(false) if err := gatewayServer.Shutdown(ctx); err != nil { log.Error().Err(err).Msg("gatewayServer.Shutdown") } else { log.Info().Msg("gatewayServer shut down correctly") } if err := statusServer.Shutdown(ctx); err != nil { log.Error().Err(err).Msg("statusServer.Shutdown") } else { log.Info().Msg("statusServer shut down correctly") } if err := metricsServer.Shutdown(ctx); err != nil { log.Error().Err(err).Msg("metricsServer.Shutdown") } else { log.Info().Msg("metricsServer shut down correctly") } grpcServer.GracefulStop() log.Info().Msgf("grpcServer shut down correctly") return nil }