WIP: first pass at sending CreateDB() through Messenger

This commit is contained in:
Travis 2017-03-14 17:44:23 -05:00
parent 7c2f6a9bae
commit 2b35df879c
6 changed files with 40 additions and 4 deletions

View file

@ -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))

View file

@ -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 {

View file

@ -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
}

View file

@ -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)
}

View file

@ -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
}

View file

@ -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)