Merge remote-tracking branch 'upstream/develop' into clearbit-notime

This commit is contained in:
Todd Gruben 2018-06-26 13:24:06 -05:00
commit f16426b50a
10 changed files with 170 additions and 105 deletions

6
Gopkg.lock generated
View file

@ -197,12 +197,6 @@
revision = "645ef00459ed84a119197bfb8d8205042c6df63d"
version = "v0.8.0"
[[projects]]
name = "github.com/rakyll/statik"
packages = ["fs"]
revision = "fd36b3595eb2ec8da4b8153b107f7ea08504899d"
version = "v0.1.1"
[[projects]]
name = "github.com/satori/go.uuid"
packages = ["."]

View file

@ -56,6 +56,7 @@ func testMessageMarshal(t *testing.T, m proto.Message) {
// Ensure that BroadcastReceiver can register a BroadcastHandler.
func TestBroadcast_BroadcastReceiver(t *testing.T) {
t.Skip("broadcast receiver")
path, err := ioutil.TempDir("", "pilosa-")
if err != nil {
panic(err)
@ -67,24 +68,24 @@ func TestBroadcast_BroadcastReceiver(t *testing.T) {
if err != nil {
t.Fatalf("setting up server: %v", err)
}
s := com.Server
// s := com.Server
sbr := NewSimpleBroadcastReceiver()
sbh := NewSimpleBroadcastHandler()
// sbr := NewSimpleBroadcastReceiver()
// sbh := NewSimpleBroadcastHandler()
s.BroadcastReceiver = sbr
s.BroadcastReceiver.Start(sbh)
// s.BroadcastReceiver = sbr
// s.BroadcastReceiver.Start(sbh)
msg := &internal.DeleteIndexMessage{
Index: "i",
}
// msg := &internal.DeleteIndexMessage{
// Index: "i",
// }
s.BroadcastReceiver.(*SimpleBroadcastReceiver).Receive(msg)
// s.BroadcastReceiver.(*SimpleBroadcastReceiver).Receive(msg)
// Make sure the message received is what was sentd
if !reflect.DeepEqual(sbh.receivedMessage, msg) {
t.Fatalf("unexpected message: %s", sbh.receivedMessage)
}
// // Make sure the message received is what was sentd
// if !reflect.DeepEqual(sbh.receivedMessage, msg) {
// t.Fatalf("unexpected message: %s", sbh.receivedMessage)
// }
}
type SimpleBroadcastReceiver struct {

View file

@ -886,11 +886,6 @@ func (c *Cluster) open() error {
return errors.Wrap(err, "adding local node")
}
// Start the EventReceiver.
if err := c.EventReceiver.Start(c); err != nil {
return fmt.Errorf("starting EventReceiver: %v", err)
}
// Open MemberSet communication.
if err := c.MemberSet.Open(c.Node); err != nil {
return fmt.Errorf("opening MemberSet: %v", err)

View file

@ -1606,6 +1606,9 @@ func (e *Executor) translateCall(index string, idx *Index, c *pql.Call) error {
}
// Translate column key.
if idx.Keys() {
if c.Args[colKey] != nil && !isString(c.Args[colKey]) {
return errors.New("column value must be a string when index 'keys' option enabled")
}
if value := callArgString(c, colKey); value != "" {
ids, err := e.TranslateStore.TranslateColumnsToUint64(index, []string{value})
if err != nil {
@ -1613,6 +1616,10 @@ func (e *Executor) translateCall(index string, idx *Index, c *pql.Call) error {
}
c.Args[colKey] = ids[0]
}
} else {
if isString(c.Args[colKey]) {
return errors.New("string 'col' value not allowed unless index 'keys' option enabled")
}
}
// Translate row key, if field is specified & key exists.
@ -1622,6 +1629,9 @@ func (e *Executor) translateCall(index string, idx *Index, c *pql.Call) error {
return ErrFieldNotFound
}
if field.Keys() {
if c.Args[rowKey] != nil && !isString(c.Args[rowKey]) {
return errors.New("row value must be a string when field 'keys' option enabled")
}
if value := callArgString(c, rowKey); value != "" {
ids, err := e.TranslateStore.TranslateRowsToUint64(index, fieldName, []string{value})
if err != nil {
@ -1629,6 +1639,10 @@ func (e *Executor) translateCall(index string, idx *Index, c *pql.Call) error {
}
c.Args[rowKey] = ids[0]
}
} else {
if isString(c.Args[rowKey]) {
return errors.New("string 'row' value not allowed unless field 'keys' option enabled")
}
}
}
@ -1801,3 +1815,8 @@ func callArgString(call *pql.Call, key string) string {
s, _ := value.(string)
return s
}
func isString(v interface{}) bool {
_, ok := v.(string)
return ok
}

View file

@ -27,6 +27,7 @@ import (
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/pql"
"github.com/pilosa/pilosa/test"
"github.com/pkg/errors"
)
// Ensure a bitmap query can be executed.
@ -264,36 +265,107 @@ func TestExecutor_Execute_Count(t *testing.T) {
}
// Ensure a set query can be executed.
func TestExecutor_Execute_Set(t *testing.T) {
hldr := test.MustOpenHolder()
defer hldr.Close()
func TestExecutor_Execute_SetBit(t *testing.T) {
t.Run("ID", func(t *testing.T) {
cmd := test.MustRunMainWithCluster(t, 1)[0]
holder := cmd.Server.Holder()
hldr := test.Holder{Holder: holder}
hldr.SetBit("i", "f", 1, 0)
// set a bit so the view gets created.
hldr.SetBit("i", "f", 1, 0)
t.Run("OK", func(t *testing.T) {
hldr.ClearBit("i", "f", 11, 1)
if n := hldr.Row("i", "f", 11).Count(); n != 0 {
t.Fatalf("unexpected bitmap count: %d", n)
}
e := test.NewExecutor(hldr.Holder, pilosa.NewTestCluster(1))
if n := hldr.Row("i", "f", 11).Count(); n != 0 {
t.Fatalf("unexpected bitmap count: %d", n)
}
if res, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(1, f=11)`}); err != nil {
t.Fatal(err)
} else {
if !res.Results[0].(bool) {
t.Fatalf("expected column changed")
}
}
if res, err := e.Execute(context.Background(), "i", test.MustParse(`Set(1, f=11)`), nil, nil); err != nil {
t.Fatal(err)
} else {
if !res[0].(bool) {
t.Fatalf("expected column changed")
}
}
if n := hldr.Row("i", "f", 11).Count(); n != 1 {
t.Fatalf("unexpected bitmap count: %d", n)
}
if res, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(1, f=11)`}); err != nil {
t.Fatal(err)
} else {
if res.Results[0].(bool) {
t.Fatalf("expected column unchanged")
}
}
})
if n := hldr.Row("i", "f", 11).Count(); n != 1 {
t.Fatalf("unexpected bitmap count: %d", n)
}
if res, err := e.Execute(context.Background(), "i", test.MustParse(`Set(1, f=11)`), nil, nil); err != nil {
t.Fatal(err)
} else {
if res[0].(bool) {
t.Fatalf("expected column unchanged")
}
}
t.Run("ErrInvalidColValueType", func(t *testing.T) {
if _, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set("foo", f=1)`}); err == nil || errors.Cause(err).Error() != `string 'col' value not allowed unless index 'keys' option enabled` {
t.Fatalf("The error is: '%v'", err)
}
})
t.Run("ErrInvalidRowValueType", func(t *testing.T) {
if _, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(2, f="bar")`}); err == nil || errors.Cause(err).Error() != `string 'row' value not allowed unless field 'keys' option enabled` {
t.Fatal(err)
}
})
})
t.Run("Keys", func(t *testing.T) {
cmd := test.MustRunMainWithCluster(t, 1)[0]
holder := cmd.Server.Holder()
hldr := test.Holder{Holder: holder}
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{Keys: true})
t.Run("OK", func(t *testing.T) {
hldr.SetBit("i", "f", 1, 0)
if n := hldr.Row("i", "f", 11).Count(); n != 0 {
t.Fatalf("unexpected bitmap count: %d", n)
}
if res, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set("foo", f=11)`}); err != nil {
t.Fatal(err)
} else {
if !res.Results[0].(bool) {
t.Fatalf("expected column changed")
}
}
if n := hldr.Row("i", "f", 11).Count(); n != 1 {
t.Fatalf("unexpected bitmap count: %d", n)
}
if res, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set("foo", f=11)`}); err != nil {
t.Fatal(err)
} else {
if res.Results[0].(bool) {
t.Fatalf("expected column unchanged")
}
}
})
t.Run("ErrInvalidColValueType", func(t *testing.T) {
if err := index.DeleteField("f"); err != nil {
t.Fatal(err)
}
if _, err := index.CreateField("f", pilosa.FieldOptions{}); err != nil {
t.Fatal(err)
}
if _, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(2, f=1)`}); err == nil || errors.Cause(err).Error() != `column value must be a string when index 'keys' option enabled` {
t.Fatal(err)
}
})
t.Run("ErrInvalidRowValueType", func(t *testing.T) {
index := hldr.MustCreateIndexIfNotExists("inokey", pilosa.IndexOptions{})
if _, err := index.CreateField("f", pilosa.FieldOptions{Keys: true}); err != nil {
t.Fatal(err)
}
if _, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "inokey", Query: `Set(2, f=1)`}); err == nil || errors.Cause(err).Error() != `row value must be a string when field 'keys' option enabled` {
t.Fatal(err)
}
})
})
}
// Ensure old PQL syntax doesn't break anything too badly.

View file

@ -33,7 +33,6 @@ import (
)
// Ensure GossipMemberSet implements interfaces.
var _ pilosa.BroadcastReceiver = &GossipMemberSet{}
var _ memberlist.Delegate = &GossipMemberSet{}
// GossipMemberSet represents a gossip implementation of MemberSet using memberlist.
@ -45,19 +44,15 @@ type GossipMemberSet struct {
broadcasts *memberlist.TransmitLimitedQueue
statusHandler pilosa.StatusHandler
config *gossipConfig
pserver *pilosa.Server
config *gossipConfig
Logger pilosa.Logger
logger *log.Logger
transport *Transport
}
// Start implements the BroadcastReceiver interface and sets the BroadcastHandler.
func (g *GossipMemberSet) Start(h pilosa.BroadcastHandler) error {
g.handler = h
return nil
gossipEventReceiver *GossipEventReceiver
}
// GetBindAddr returns the gossip bind address based on config and auto bind port.
@ -69,13 +64,16 @@ func (g *GossipMemberSet) GetBindAddr() string {
// Open implements the MemberSet interface to start network activity.
func (g *GossipMemberSet) Open(n *pilosa.Node) error {
err := g.gossipEventReceiver.Start(g.pserver)
if err != nil {
return errors.Wrap(err, "starting event delegate")
}
if g.handler == nil {
return fmt.Errorf("must call Start(pilosa.BroadcastHandler) before calling Open()")
}
g.node = n
err := error(nil)
g.mu.Lock()
g.memberlist, err = memberlist.Create(g.config.memberlistConfig)
g.mu.Unlock()
@ -166,7 +164,7 @@ func WithLogger(logger *log.Logger) GossipMemberSetOption {
}
// NewGossipMemberSet returns a new instance of GossipMemberSet based on options.
func NewGossipMemberSet(name string, host string, cfg Config, ger *GossipEventReceiver, sh pilosa.StatusHandler, options ...GossipMemberSetOption) (*GossipMemberSet, error) {
func NewGossipMemberSet(name string, host string, cfg Config, s *pilosa.Server, options ...GossipMemberSetOption) (*GossipMemberSet, error) {
g := &GossipMemberSet{
Logger: pilosa.NopLogger,
}
@ -177,6 +175,10 @@ func NewGossipMemberSet(name string, host string, cfg Config, ger *GossipEventRe
return nil, errors.Wrap(err, "executing option")
}
}
ger := NewGossipEventReceiver(g.logger)
g.gossipEventReceiver = ger
g.handler = s
if g.transport == nil {
port, err := strconv.Atoi(cfg.Port)
@ -232,7 +234,7 @@ func NewGossipMemberSet(name string, host string, cfg Config, ger *GossipEventRe
gossipSeeds: cfg.Seeds,
}
g.statusHandler = sh
g.pserver = s
return g, nil
}
@ -270,7 +272,7 @@ func (g *GossipMemberSet) GetBroadcasts(overhead, limit int) [][]byte {
// LocalState implementation of the memberlist.Delegate interface
// sends this Node's state data.
func (g *GossipMemberSet) LocalState(join bool) []byte {
pb, err := g.statusHandler.LocalStatus()
pb, err := g.pserver.LocalStatus()
if err != nil {
g.Logger.Printf("error getting local state, err=%s", err)
return []byte{}
@ -294,7 +296,7 @@ func (g *GossipMemberSet) MergeRemoteState(buf []byte, join bool) {
g.Logger.Printf("error unmarshalling nodestate data, err=%s", err)
return
}
err := g.statusHandler.HandleRemoteStatus(&pb)
err := g.pserver.HandleRemoteStatus(&pb)
if err != nil {
g.Logger.Printf("merge state error: %s", err)
}
@ -309,11 +311,11 @@ type GossipEventReceiver struct {
ch chan memberlist.NodeEvent
eventHandler pilosa.EventHandler
Logger pilosa.Logger
Logger *log.Logger
}
// NewGossipEventReceiver returns a new instance of GossipEventReceiver.
func NewGossipEventReceiver(logger pilosa.Logger) *GossipEventReceiver {
func NewGossipEventReceiver(logger *log.Logger) *GossipEventReceiver {
return &GossipEventReceiver{
ch: make(chan memberlist.NodeEvent, 1),
Logger: logger,

View file

@ -61,10 +61,9 @@ type Server struct {
clusterDisabled bool
// External
BroadcastReceiver BroadcastReceiver
systemInfo SystemInfo
gcNotifier GCNotifier
logger Logger
systemInfo SystemInfo
gcNotifier GCNotifier
logger Logger
NodeID string
URI URI
@ -207,12 +206,11 @@ func OptServerClusterDisabled(disabled bool, hosts []string) ServerOption {
// NewServer returns a new instance of Server.
func NewServer(opts ...ServerOption) (*Server, error) {
s := &Server{
closing: make(chan struct{}),
Cluster: NewCluster(),
holder: NewHolder(),
BroadcastReceiver: NopBroadcastReceiver,
diagnostics: NewDiagnosticsCollector(DefaultDiagnosticServer),
systemInfo: NewNopSystemInfo(),
closing: make(chan struct{}),
Cluster: NewCluster(),
holder: NewHolder(),
diagnostics: NewDiagnosticsCollector(DefaultDiagnosticServer),
systemInfo: NewNopSystemInfo(),
gcNotifier: NopGCNotifier,
@ -246,11 +244,8 @@ func NewServer(opts ...ServerOption) (*Server, error) {
// Initialize translation database.
s.translateFile = NewTranslateFile()
s.translateFile.Path = filepath.Join(path, "keys")
s.translateFile.Path = filepath.Join(path, ".keys")
s.translateFile.PrimaryTranslateStore = s.primaryTranslateStore
if err := s.translateFile.Open(); err != nil {
return nil, err
}
// Get or create NodeID.
s.NodeID = s.LoadNodeID()
@ -290,6 +285,11 @@ func (s *Server) Open() error {
log.Println(errors.Wrap(err, "logging startup"))
}
// Initialize id-key storage.
if err := s.translateFile.Open(); err != nil {
return err
}
// Cluster settings.
s.Cluster.Broadcaster = s
s.Cluster.MaxWritesPerRequest = s.maxWritesPerRequest
@ -297,11 +297,6 @@ func (s *Server) Open() error {
// Initialize Holder.
s.holder.Broadcaster = s
// Start the BroadcastReceiver.
if err := s.BroadcastReceiver.Start(s); err != nil {
return fmt.Errorf("starting BroadcastReceiver: %v", err)
}
// Open Cluster management.
if err := s.Cluster.open(); err != nil {
return fmt.Errorf("opening Cluster: %v", err)
@ -711,6 +706,11 @@ func (s *Server) monitorRuntime() {
}
}
// ReceiveEvent implements the EventHandler interface.
func (s *Server) ReceiveEvent(e *NodeEvent) error {
return s.Cluster.ReceiveEvent(e)
}
// countOpenFiles on operating systems that support lsof.
func countOpenFiles() (int, error) {
switch runtime.GOOS {

View file

@ -75,6 +75,7 @@ type Command struct {
logger loggerLogger
Handler pilosa.Handler
API *pilosa.API
ln net.Listener
serverOptions []pilosa.ServerOption
@ -274,14 +275,14 @@ func (m *Command) SetupServer() error {
return errors.Wrap(err, "new server")
}
api, err := pilosa.NewAPI(pilosa.OptAPIServer(m.Server))
m.API, err = pilosa.NewAPI(pilosa.OptAPIServer(m.Server))
if err != nil {
return errors.Wrap(err, "new api")
}
m.Handler, err = http.NewHandler(
http.OptHandlerAllowedOrigins(m.Config.Handler.AllowedOrigins),
http.OptHandlerAPI(api),
http.OptHandlerAPI(m.API),
http.OptHandlerLogger(m.logger),
http.OptHandlerListener(m.ln),
)
@ -318,13 +319,10 @@ func (m *Command) SetupNetworking() error {
m.Server.Cluster.Node.IsCoordinator = true
}
gossipEventReceiver := gossip.NewGossipEventReceiver(m.logger)
m.Server.Cluster.EventReceiver = gossipEventReceiver
gossipMemberSet, err := gossip.NewGossipMemberSet(
m.Server.NodeID,
m.Server.URI.Host(),
m.Config.Gossip,
gossipEventReceiver,
m.Server,
gossip.WithLogger(m.logger.Logger()),
gossip.WithTransport(transport),
@ -332,9 +330,7 @@ func (m *Command) SetupNetworking() error {
if err != nil {
return errors.Wrap(err, "getting memberset")
}
gossipMemberSet.Logger = m.logger
m.Server.Cluster.MemberSet = gossipMemberSet
m.Server.BroadcastReceiver = gossipMemberSet
return nil
}

File diff suppressed because one or more lines are too long

View file

@ -224,10 +224,6 @@ func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string) (
return seed, err
}
if err = m.Server.BroadcastReceiver.Start(m.Server); err != nil {
return seed, err
}
m.Server.Cluster.Static = false
go func() {