From 2b35df879cb3c7e0b474ac9c47116db789fabe8e Mon Sep 17 00:00:00 2001 From: Travis Date: Tue, 14 Mar 2017 17:44:23 -0500 Subject: [PATCH] WIP: first pass at sending `CreateDB()` through Messenger --- client.go | 1 + cluster.go | 4 ++++ index.go | 17 ++++++++++++++++- messenger.go | 5 +++++ server.go | 8 +++++--- server/server.go | 9 +++++++++ 6 files changed, 40 insertions(+), 4 deletions(-) diff --git a/client.go b/client.go index e4401a69a..d340ce640 100644 --- a/client.go +++ b/client.go @@ -159,6 +159,7 @@ func (c *Client) CreateDB(ctx context.Context, db string, opt DBOptions) error { case http.StatusOK: return nil // ok case http.StatusConflict: + fmt.Println("ErrDatabaseExists: 1") return ErrDatabaseExists default: return errors.New(string(body)) diff --git a/cluster.go b/cluster.go index 8268fe508..061f76c62 100644 --- a/cluster.go +++ b/cluster.go @@ -285,6 +285,7 @@ func (h *HTTPNodeSet) ReceiveMessage(pb proto.Message) error { } func (h *HTTPNodeSet) sendNodeMessage(node *Node, msg []byte) error { + fmt.Println("sendNodeMessage:", node.Host) var client *http.Client client = http.DefaultClient @@ -302,7 +303,9 @@ func (h *HTTPNodeSet) sendNodeMessage(node *Node, msg []byte) error { req.Header.Set("Content-Type", "application/x-protobuf") // Send request to remote node. + fmt.Println("Send") resp, err := client.Do(req) + fmt.Println("Got back") if err != nil { return err } @@ -314,6 +317,7 @@ func (h *HTTPNodeSet) sendNodeMessage(node *Node, msg []byte) error { if err != nil { return err } + fmt.Println("code:", resp.StatusCode) // Check status code. if resp.StatusCode != http.StatusOK { diff --git a/index.go b/index.go index 01409a277..f7269d670 100644 --- a/index.go +++ b/index.go @@ -187,12 +187,14 @@ func (i *Index) DBs() []*DB { } // CreateDB creates a database. +// An error is returned if the database already exists. func (i *Index) CreateDB(name string, opt DBOptions) (*DB, error) { i.mu.Lock() defer i.mu.Unlock() // Ensure db doesn't already exist. if i.dbs[name] != nil { + fmt.Println("ErrDatabaseExists: 2") return nil, ErrDatabaseExists } return i.createDB(name, opt) @@ -204,7 +206,7 @@ func (i *Index) CreateDBIfNotExists(name string, opt DBOptions) (*DB, error) { i.mu.Lock() defer i.mu.Unlock() - // Find frame in cache first. + // Find database in cache first. if db := i.dbs[name]; db != nil { return db, nil } @@ -239,6 +241,13 @@ func (i *Index) createDB(name string, opt DBOptions) (*DB, error) { i.Stats.Count("dbN", 1) + // Send a CreateDB message + i.Messenger.SendMessage( + &internal.CreateDBMessage{ + DB: name, + ColumnLabel: opt.ColumnLabel, + }) + return db, nil } @@ -360,6 +369,12 @@ func (i *Index) HandleMessage(pb proto.Message) error { if err != nil { return err } + case *internal.CreateDBMessage: + opt := DBOptions{ColumnLabel: obj.ColumnLabel} + _, err := i.CreateDB(obj.DB, opt) + if err != nil { + return err + } } return nil } diff --git a/messenger.go b/messenger.go index 115627ae4..5e2d128e2 100644 --- a/messenger.go +++ b/messenger.go @@ -34,6 +34,7 @@ type Messenger interface { const ( MessageTypeCreateSlice = 1 MessageTypeDeleteDB = 2 + MessageTypeCreateDB = 3 ) func MarshalMessage(m proto.Message) ([]byte, error) { @@ -43,6 +44,8 @@ func MarshalMessage(m proto.Message) ([]byte, error) { typ = MessageTypeCreateSlice case *internal.DeleteDBMessage: typ = MessageTypeDeleteDB + case *internal.CreateDBMessage: + typ = MessageTypeCreateDB default: return nil, fmt.Errorf("message type not implemented for marshalling: %s", reflect.TypeOf(obj)) } @@ -62,6 +65,8 @@ func UnmarshalMessage(buf []byte) (proto.Message, error) { m = &internal.CreateSliceMessage{} case MessageTypeDeleteDB: m = &internal.DeleteDBMessage{} + case MessageTypeCreateDB: + m = &internal.CreateDBMessage{} default: return nil, fmt.Errorf("invalid message type: %d", typ) } diff --git a/server.go b/server.go index f4cef42f6..cde64c3ca 100644 --- a/server.go +++ b/server.go @@ -120,9 +120,11 @@ func (s *Server) Open() error { go func() { http.Serve(ln, s.Handler) }() // Start background monitoring. - s.wg.Add(2) - go func() { defer s.wg.Done(); s.monitorAntiEntropy() }() - go func() { defer s.wg.Done(); s.monitorMaxSlices() }() + /* + s.wg.Add(2) + go func() { defer s.wg.Done(); s.monitorAntiEntropy() }() + go func() { defer s.wg.Done(); s.monitorMaxSlices() }() + */ return nil } diff --git a/server/server.go b/server/server.go index 251a054aa..da83da52c 100644 --- a/server/server.go +++ b/server/server.go @@ -94,6 +94,15 @@ func (m *Command) Run(args ...string) (err error) { } m.Server.Cluster = m.Config.PilosaCluster() + // Setup Messenger. + fmt.Fprintf(m.Stderr, "Using Messenger type: %s\n", m.Config.Cluster.MessengerType) + m.Server.Messenger = m.Server.Cluster.NodeSet.(pilosa.Messenger) + m.Server.Handler.Messenger = m.Server.Messenger + m.Server.Index.Messenger = m.Server.Messenger + + // Set message handler. + m.Server.Cluster.NodeSet.SetMessageHandler(m.Server.Index.HandleMessage) + // Set configuration options. m.Server.AntiEntropyInterval = time.Duration(m.Config.AntiEntropy.Interval)