From 26852f1850c0819b27a503637d900a3c74ac4bc1 Mon Sep 17 00:00:00 2001 From: Willy Kloucek Date: Fri, 18 Feb 2022 11:51:56 +0100 Subject: [PATCH] make nats killable --- nats/pkg/command/server.go | 1 + nats/pkg/server/nats/nats.go | 16 +++++++++++++--- 2 files changed, 14 insertions(+), 3 deletions(-) diff --git a/nats/pkg/command/server.go b/nats/pkg/command/server.go index be2e0ff46..c7e613aca 100644 --- a/nats/pkg/command/server.go +++ b/nats/pkg/command/server.go @@ -36,6 +36,7 @@ func Server(cfg *config.Config) *cli.Command { defer cancel() natsServer, err := nats.NewNATSServer( + ctx, nats.Host(cfg.Nats.Host), nats.Port(cfg.Nats.Port), nats.Logger(logging.NewLogWrapper(logger)), diff --git a/nats/pkg/server/nats/nats.go b/nats/pkg/server/nats/nats.go index 4defe9eed..28d512fbb 100644 --- a/nats/pkg/server/nats/nats.go +++ b/nats/pkg/server/nats/nats.go @@ -1,6 +1,7 @@ package nats import ( + "context" "time" natsServer "github.com/nats-io/nats-server/v2/server" @@ -8,14 +9,18 @@ import ( ) type NATSServer struct { + ctx context.Context + natsOpts *natsServer.Options stanOpts *stanServer.Options server *stanServer.StanServer } -func NewNATSServer(opts ...Option) (*NATSServer, error) { +func NewNATSServer(ctx context.Context, opts ...Option) (*NATSServer, error) { + server := &NATSServer{ + ctx: ctx, natsOpts: &stanServer.DefaultNatsServerOptions, stanOpts: stanServer.GetDefaultOptions(), } @@ -28,7 +33,6 @@ func NewNATSServer(opts ...Option) (*NATSServer, error) { } func (n *NATSServer) ListenAndServe() (err error) { - n.server, err = stanServer.RunServerWithOpts( n.stanOpts, n.natsOpts, @@ -37,15 +41,21 @@ func (n *NATSServer) ListenAndServe() (err error) { return err } + defer n.Shutdown() + 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 + // check if the NATs server is still running if n.server.State() == stanServer.Shutdown { return nil } + // check if context was cancelled + if n.ctx.Err() != nil { + return nil + } time.Sleep(1 * time.Second) } }