diff --git a/gossip/gossip.go b/gossip/gossip.go index 2b983376f..a5b4e9299 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -93,6 +93,17 @@ func (g *memberSet) Open() (err error) { return nil } +// Close attempts to gracefully leave the cluster, and finally calls shutdown +// after (at most) a timeout period. +func (g *memberSet) Close() error { + leaveErr := g.memberlist.Leave(5 * time.Second) + shutdownErr := g.memberlist.Shutdown() + if leaveErr != nil || shutdownErr != nil { + return fmt.Errorf("leaving: '%v', shutting down: '%v'", leaveErr, shutdownErr) + } + return nil +} + // joinWithRetry wraps the standard memberlist Join function in a retry. func (g *memberSet) joinWithRetry(hosts []string) error { err := retry(60, 2*time.Second, func() error { diff --git a/http/handler.go b/http/handler.go index 4946f69e0..d85d3a8a5 100644 --- a/http/handler.go +++ b/http/handler.go @@ -55,6 +55,8 @@ type Handler struct { ln net.Listener + closeTimeout time.Duration + server *http.Server } @@ -109,10 +111,20 @@ func OptHandlerListener(ln net.Listener) handlerOption { } } +// OptHandlerCloseTimeout controls how long to wait for the http Server to +// shutdown cleanly before forcibly destroying it. Default is 30 seconds. +func OptHandlerCloseTimeout(d time.Duration) handlerOption { + return func(h *Handler) error { + h.closeTimeout = d + return nil + } +} + // NewHandler returns a new instance of Handler with a default logger. func NewHandler(opts ...handlerOption) (*Handler, error) { handler := &Handler{ - logger: pilosa.NopLogger, + logger: pilosa.NopLogger, + closeTimeout: time.Second * 30, } handler.Handler = newRouter(handler) handler.populateValidators() @@ -146,10 +158,16 @@ func (h *Handler) Serve() error { return nil } +// Close tries to cleanly shutdown the HTTP server, and failing that, after a +// timeout, calls Server.Close. func (h *Handler) Close() error { - // TODO: timeout? - err := h.server.Shutdown(context.Background()) - return errors.Wrap(err, "shutdown http server") + deadlineCtx, cancelFunc := context.WithDeadline(context.Background(), time.Now().Add(h.closeTimeout)) + defer cancelFunc() + err := h.server.Shutdown(deadlineCtx) + if err != nil { + err = h.server.Close() + } + return errors.Wrap(err, "shutdown/close http server") } func (h *Handler) populateValidators() { diff --git a/server.go b/server.go index 91c2d813a..db577fb86 100644 --- a/server.go +++ b/server.go @@ -369,17 +369,28 @@ func (s *Server) Close() error { close(s.closing) s.wg.Wait() + var errh error + var errt error + var errc error if s.cluster != nil { - s.cluster.close() + errc = s.cluster.close() } if s.holder != nil { - s.holder.Close() + errh = s.holder.Close() } if s.translateFile != nil { - s.translateFile.Close() + errt = s.translateFile.Close() } - - return nil + // prefer to return holder error over translateFile error over cluster + // error. This order is somewhat arbitrary. It would be better if we had + // some way to combine all the errors, but probably not important enough to + // warrant the extra complexity. + if errh != nil { + return errors.Wrap(errh, "closing holder") + } else if errt != nil { + return errors.Wrap(errt, "closing translateFile") + } + return errors.Wrap(errc, "closing cluster") } // loadNodeID gets NodeID from disk, or creates a new value. diff --git a/server/server.go b/server/server.go index e83cc1c1c..9060256af 100644 --- a/server/server.go +++ b/server/server.go @@ -20,7 +20,7 @@ package server import ( - "fmt" + "crypto/tls" "io" "log" "math/rand" @@ -31,7 +31,7 @@ import ( "syscall" "time" - "crypto/tls" + "golang.org/x/sync/errgroup" "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/boltdb" @@ -62,6 +62,7 @@ type Command struct { // Gossip transport gossipTransport *gossip.Transport + gossipMemberSet io.Closer // Standard input/output *pilosa.CmdIO @@ -75,9 +76,10 @@ type Command struct { logOutput io.Writer logger loggerLogger - Handler pilosa.Handler - API *pilosa.API - ln net.Listener + Handler pilosa.Handler + API *pilosa.API + ln net.Listener + closeTimeout time.Duration serverOptions []pilosa.ServerOption } @@ -91,6 +93,13 @@ func OptCommandServerOptions(opts ...pilosa.ServerOption) CommandOption { } } +func OptCommandCloseTimeout(d time.Duration) CommandOption { + return func(c *Command) error { + c.closeTimeout = d + return nil + } +} + // NewCommand returns a new instance of Main. func NewCommand(stdin io.Reader, stdout, stderr io.Writer, opts ...CommandOption) *Command { c := &Command{ @@ -294,6 +303,7 @@ func (m *Command) SetupServer() error { http.OptHandlerAPI(m.API), http.OptHandlerLogger(m.logger), http.OptHandlerListener(m.ln), + http.OptHandlerCloseTimeout(m.closeTimeout), ) return errors.Wrap(err, "new handler") @@ -326,6 +336,8 @@ func (m *Command) setupNetworking() error { if err != nil { return errors.Wrap(err, "getting memberset") } + m.gossipMemberSet = gossipMemberSet + return errors.Wrap(gossipMemberSet.Open(), "opening gossip memberset") } @@ -338,17 +350,18 @@ func (m *Command) GossipTransport() *gossip.Transport { // Close shuts down the server. func (m *Command) Close() error { - var logErr error - handlerErr := m.Handler.Close() - serveErr := m.Server.Close() + defer close(m.done) + eg := errgroup.Group{} + eg.Go(m.Handler.Close) + eg.Go(m.Server.Close) + if m.gossipMemberSet != nil { + eg.Go(m.gossipMemberSet.Close) + } if closer, ok := m.logOutput.(io.Closer); ok { - logErr = closer.Close() + eg.Go(closer.Close) } - close(m.done) - if serveErr != nil || logErr != nil || handlerErr != nil { - return fmt.Errorf("closing server: '%v', closing logs: '%v', closing handler: '%v'", serveErr, logErr, handlerErr) - } - return nil + err := eg.Wait() + return errors.Wrap(err, "closing everything") } // newStatsClient creates a stats client from the config diff --git a/test/pilosa.go b/test/pilosa.go index e1022ce8b..001949906 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -23,6 +23,7 @@ import ( "os" "strings" "testing" + "time" "github.com/pilosa/pilosa/http" "github.com/pilosa/pilosa/server" @@ -55,6 +56,11 @@ func newCommand(opts ...server.CommandOption) *Command { panic(err) } + // set aggressive close timeout by default to avoid hanging tests. This was + // a problem with PDK tests which used go-pilosa as well. We put it at the + // beginning of the option slice so that it can be overridden by user-passed + // options. + opts = append([]server.CommandOption{server.OptCommandCloseTimeout(time.Millisecond * 2)}, opts...) m := &Command{Command: server.NewCommand(os.Stdin, os.Stdout, os.Stderr, opts...), commandOptions: opts} m.Config.DataDir = path m.Config.Bind = "http://localhost:0" diff --git a/translate_mapsize_386.go b/translate_mapsize_386.go index b8beafa13..259b735ac 100644 --- a/translate_mapsize_386.go +++ b/translate_mapsize_386.go @@ -1,5 +1,5 @@ package pilosa -// DefaultMapSize is the default size of mapped memory for the translate store. +// defaultMapSize is the default size of mapped memory for the translate store. // It is passed as an int to syscall.Mmap and so must be < 2^31 -const DefaultMapSize = (1 << 31) - 1 // 2GB +const defaultMapSize = (1 << 31) - 1 // 2GB