mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-09 22:51:02 +00:00
Merge branch 'master' into time-range
This commit is contained in:
commit
493e77ae41
6 changed files with 84 additions and 25 deletions
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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() {
|
||||
|
|
|
|||
21
server.go
21
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.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue