From 3a60056bb3c45325d283309ee300f307068cb1ec Mon Sep 17 00:00:00 2001 From: Willy Kloucek Date: Fri, 18 Feb 2022 11:44:26 +0100 Subject: [PATCH] transfer logic into nats server package --- nats/pkg/command/server.go | 50 +++++++------------------------ nats/pkg/server/nats/nats.go | 52 ++++++++++++++++++++++++++++----- nats/pkg/server/nats/options.go | 16 ++++------ 3 files changed, 61 insertions(+), 57 deletions(-) diff --git a/nats/pkg/command/server.go b/nats/pkg/command/server.go index f12538ab6..be2e0ff46 100644 --- a/nats/pkg/command/server.go +++ b/nats/pkg/command/server.go @@ -3,7 +3,6 @@ package command import ( "context" "fmt" - "time" "github.com/oklog/run" @@ -12,9 +11,6 @@ import ( "github.com/owncloud/ocis/nats/pkg/logging" "github.com/owncloud/ocis/nats/pkg/server/nats" "github.com/urfave/cli/v2" - - // TODO: .Logger Option on events/server would make this import redundant - stanServer "github.com/nats-io/nats-streaming-server/server" ) // Server is the entrypoint for the server command. @@ -39,46 +35,22 @@ func Server(cfg *config.Config) *cli.Command { defer cancel() - var natsServer *stanServer.StanServer + natsServer, err := nats.NewNATSServer( + nats.Host(cfg.Nats.Host), + nats.Port(cfg.Nats.Port), + nats.Logger(logging.NewLogWrapper(logger)), + ) + if err != nil { + return err + } gr.Add(func() error { - var err error - - natsServer, err = nats.RunNatsServer( - nats.Host(cfg.Nats.Host), - nats.Port(cfg.Nats.Port), - nats.StanOpts( - func(o *stanServer.Options) { - o.CustomLogger = logging.NewLogWrapper(logger) - }, - ), - ) - - if err != nil { - return err - } - - errChan := make(chan error) - - go func() { - for { - // check if NATs server has an encountered an error - if err := natsServer.LastError(); err != nil { - errChan <- err - return - } - if ctx.Err() != nil { - return // context closed - } - time.Sleep(1 * time.Second) - } - }() - + err := make(chan error) select { case <-ctx.Done(): return nil - case err = <-errChan: - return err + case err <- natsServer.ListenAndServe(): + return <-err } }, func(_ error) { diff --git a/nats/pkg/server/nats/nats.go b/nats/pkg/server/nats/nats.go index 398a4c966..4defe9eed 100644 --- a/nats/pkg/server/nats/nats.go +++ b/nats/pkg/server/nats/nats.go @@ -1,17 +1,55 @@ package nats import ( + "time" + + natsServer "github.com/nats-io/nats-server/v2/server" stanServer "github.com/nats-io/nats-streaming-server/server" ) -// RunNatsServer runs the nats streaming server -func RunNatsServer(opts ...Option) (*stanServer.StanServer, error) { - natsOpts := stanServer.DefaultNatsServerOptions - stanOpts := stanServer.GetDefaultOptions() +type NATSServer struct { + natsOpts *natsServer.Options + stanOpts *stanServer.Options + + server *stanServer.StanServer +} + +func NewNATSServer(opts ...Option) (*NATSServer, error) { + server := &NATSServer{ + natsOpts: &stanServer.DefaultNatsServerOptions, + stanOpts: stanServer.GetDefaultOptions(), + } for _, o := range opts { - o(&natsOpts, stanOpts) + o(server.natsOpts, server.stanOpts) } - s, err := stanServer.RunServerWithOpts(stanOpts, &natsOpts) - return s, err + + return server, nil +} + +func (n *NATSServer) ListenAndServe() (err error) { + + n.server, err = stanServer.RunServerWithOpts( + n.stanOpts, + n.natsOpts, + ) + if err != nil { + return err + } + + for { + // check if NATs server has an encountered an error + if err := n.server.LastError(); err != nil { + return err + } + // check if th NATs server is still running + if n.server.State() == stanServer.Shutdown { + return nil + } + time.Sleep(1 * time.Second) + } +} + +func (n *NATSServer) Shutdown() { + n.server.Shutdown() } diff --git a/nats/pkg/server/nats/options.go b/nats/pkg/server/nats/options.go index 74e54e51f..0189141ca 100644 --- a/nats/pkg/server/nats/options.go +++ b/nats/pkg/server/nats/options.go @@ -2,6 +2,7 @@ package nats import ( natsServer "github.com/nats-io/nats-server/v2/server" + "github.com/nats-io/nats-streaming-server/logger" stanServer "github.com/nats-io/nats-streaming-server/server" ) @@ -22,16 +23,9 @@ func Port(port int) Option { } } -// NatsOpts allows setting Options from nats package directly -func NatsOpts(opt func(*natsServer.Options)) Option { - return func(no *natsServer.Options, _ *stanServer.Options) { - opt(no) - } -} - -// StanOpts allows setting Options from stan package directly -func StanOpts(opt func(*stanServer.Options)) Option { - return func(_ *natsServer.Options, so *stanServer.Options) { - opt(so) +// Port sets the host URL for the nats server +func Logger(logger logger.Logger) Option { + return func(no *natsServer.Options, so *stanServer.Options) { + so.CustomLogger = logger } }