Merge branch 'develop' into newpql

This commit is contained in:
Matt Jaffee 2018-06-22 09:11:00 -05:00
commit 80f5b13a97
No known key found for this signature in database
GPG key ID: 08A3DFFF987B11BF
11 changed files with 205 additions and 87 deletions

View file

@ -1010,6 +1010,9 @@ func (c *Cluster) handleNodeAction(nodeAction nodeAction) error {
func (c *Cluster) setStateAndBroadcast(state string) error {
c.SetState(state)
if c.Static {
return nil
}
// Broadcast cluster status changes to the cluster.
c.Logger.Printf("broadcasting ClusterStatus: %s", state)
return c.Broadcaster.SendSync(c.Status())

View file

@ -183,6 +183,7 @@ func TestImportCommand_InvalidFile(t *testing.T) {
// MustNewHTTPRequest creates a new HTTP request. Panic on error.
func MustNewHTTPRequest(method, urlStr string, body io.Reader) *http.Request {
req, err := http.NewRequest(method, urlStr, body)
req.Header.Add("Accept", "application/json")
if err != nil {
panic(err)
}

View file

@ -222,7 +222,7 @@ func (d *DiagnosticsCollector) EnrichWithSchemaProperties() {
bsiFieldCount := 0
timeQuantumEnabled := false
for _, index := range d.server.Holder.Indexes() {
for _, index := range d.server.holder.Indexes() {
numSlices += index.MaxSlice() + 1
numIndexes += 1
for _, field := range index.Fields() {

View file

@ -90,6 +90,7 @@ func (c *InternalClient) maxSliceByIndex(ctx context.Context) (map[string]uint64
}
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.HTTPClient.Do(req.WithContext(ctx))
@ -120,6 +121,7 @@ func (c *InternalClient) Schema(ctx context.Context) ([]*pilosa.IndexInfo, error
}
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.HTTPClient.Do(req.WithContext(ctx))
@ -195,6 +197,7 @@ func (c *InternalClient) FragmentNodes(ctx context.Context, index string, slice
}
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.HTTPClient.Do(req.WithContext(ctx))
@ -685,6 +688,7 @@ func (c *InternalClient) FragmentBlocks(ctx context.Context, uri *pilosa.URI, in
}
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.HTTPClient.Do(req.WithContext(ctx))
@ -777,6 +781,7 @@ func (c *InternalClient) ColumnAttrDiff(ctx context.Context, uri *pilosa.URI, in
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.HTTPClient.Do(req.WithContext(ctx))
@ -820,6 +825,7 @@ func (c *InternalClient) RowAttrDiff(ctx context.Context, uri *pilosa.URI, index
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.HTTPClient.Do(req.WithContext(ctx))
@ -859,6 +865,7 @@ func (c *InternalClient) SendMessage(ctx context.Context, uri *pilosa.URI, pb pr
}
req.Header.Set("Content-Type", "application/x-protobuf")
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.HTTPClient.Do(req.WithContext(ctx))

View file

@ -164,11 +164,12 @@ func (h *Handler) queryArgValidator(next http.Handler) http.Handler {
func NewRouter(handler *Handler) *mux.Router {
router := mux.NewRouter()
router.HandleFunc("/", handler.handleHome).Methods("GET")
router.HandleFunc("/cluster/message", handler.handlePostClusterMessage).Methods("POST")
router.HandleFunc("/cluster/resize/set-coordinator", handler.handlePostClusterResizeSetCoordinator).Methods("POST")
router.PathPrefix("/debug/pprof/").Handler(http.DefaultServeMux).Methods("GET")
router.Handle("/debug/vars", expvar.Handler()).Methods("GET")
router.HandleFunc("/schema", handler.handleGetSchema).Methods("GET")
router.HandleFunc("/cluster/message", handler.handlePostClusterMessage).Methods("POST")
router.HandleFunc("/cluster/resize/set-coordinator", handler.handlePostClusterResizeSetCoordinator).Methods("POST")
router.HandleFunc("/slices/max", handler.handleGetSlicesMax).Methods("GET") // TODO: deprecate, but it's being used by the client
router.HandleFunc("/status", handler.handleGetStatus).Methods("GET")
router.HandleFunc("/info", handler.handleGetInfo).Methods("GET")
@ -257,8 +258,28 @@ func (h *Handler) handleHome(w http.ResponseWriter, r *http.Request) {
http.Error(w, "Welcome. Pilosa is running. Visit https://www.pilosa.com/docs/ for more information.", http.StatusNotFound)
}
func checkHeaderAcceptJSON(header http.Header) bool {
v, found := header["Accept"]
sendError := false
if found {
sendError = true
for _, v := range v {
if v == "application/json" {
sendError = false
}
}
}
return sendError
}
// handleGetSchema handles GET /schema requests.
func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
schema := h.API.Schema(r.Context())
if err := json.NewEncoder(w).Encode(getSchemaResponse{
Indexes: schema,
@ -269,6 +290,10 @@ func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) {
// handleGetStatus handles GET /status requests.
func (h *Handler) handleGetStatus(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
status := getStatusResponse{
State: h.API.State(),
Nodes: h.API.Hosts(r.Context()),
@ -280,6 +305,10 @@ func (h *Handler) handleGetStatus(w http.ResponseWriter, r *http.Request) {
}
func (h *Handler) handleGetInfo(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
info := h.API.Info()
if err := json.NewEncoder(w).Encode(info); err != nil {
h.Logger.Printf("write info response error: %s", err)
@ -333,6 +362,10 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) {
// handleGetSlicesMax handles GET /schema requests.
func (h *Handler) handleGetSlicesMax(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
if err := json.NewEncoder(w).Encode(getSlicesMaxResponse{
Standard: h.API.MaxSlices(r.Context()),
}); err != nil {
@ -351,6 +384,10 @@ func (h *Handler) handleGetIndexes(w http.ResponseWriter, r *http.Request) {
// handleGetIndex handles GET /index/<indexname> requests.
func (h *Handler) handleGetIndex(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
indexName := mux.Vars(r)["index"]
for _, idx := range h.API.Schema(r.Context()) {
if idx.Name == indexName {
@ -429,6 +466,10 @@ type postIndexResponse struct{}
// handleDeleteIndex handles DELETE /index request.
func (h *Handler) handleDeleteIndex(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
indexName := mux.Vars(r)["index"]
err := h.API.DeleteIndex(r.Context(), indexName)
if err != nil {
@ -447,6 +488,10 @@ type deleteIndexResponse struct{}
// handlePostIndex handles POST /index request.
func (h *Handler) handlePostIndex(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
indexName := mux.Vars(r)["index"]
// Decode request.
@ -477,6 +522,10 @@ func (h *Handler) handlePostIndex(w http.ResponseWriter, r *http.Request) {
// handlePostIndexAttrDiff handles POST /index/attr/diff requests.
func (h *Handler) handlePostIndexAttrDiff(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
indexName := mux.Vars(r)["index"]
// Decode request.
@ -514,6 +563,10 @@ type postIndexAttrDiffResponse struct {
// handlePostField handles POST /field request.
func (h *Handler) handlePostField(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
indexName := mux.Vars(r)["index"]
fieldName := mux.Vars(r)["field"]
@ -592,6 +645,11 @@ type postFieldResponse struct{}
// handleDeleteField handles DELETE /field request.
func (h *Handler) handleDeleteField(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
indexName := mux.Vars(r)["index"]
fieldName := mux.Vars(r)["field"]
@ -617,6 +675,10 @@ type deleteFieldResponse struct{}
// handlePostFieldAttrDiff handles POST /field/attr/diff requests.
func (h *Handler) handlePostFieldAttrDiff(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
indexName := mux.Vars(r)["index"]
fieldName := mux.Vars(r)["field"]
@ -709,7 +771,7 @@ func (h *Handler) readURLQueryRequest(r *http.Request) (*pilosa.QueryRequest, er
// writeQueryResponse writes the response from the executor to w.
func (h *Handler) writeQueryResponse(w http.ResponseWriter, r *http.Request, resp *pilosa.QueryResponse) error {
if strings.Contains(r.Header.Get("Accept"), "application/x-protobuf") {
if checkHeaderAcceptJSON(r.Header) {
return h.writeProtobufQueryResponse(w, resp)
}
return h.writeJSONQueryResponse(w, resp)
@ -872,6 +934,10 @@ func (h *Handler) handleGetExportCSV(w http.ResponseWriter, r *http.Request) {
// handleGetFragmentNodes handles /fragment/nodes requests.
func (h *Handler) handleGetFragmentNodes(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
q := r.URL.Query()
index := q.Get("index")
@ -917,6 +983,10 @@ func (h *Handler) handleGetFragmentBlockData(w http.ResponseWriter, r *http.Requ
// handleGetFragmentBlocks handles GET /fragment/blocks requests.
func (h *Handler) handleGetFragmentBlocks(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
// Read slice parameter.
q := r.URL.Query()
slice, err := strconv.ParseUint(q.Get("slice"), 10, 64)
@ -949,6 +1019,10 @@ type getFragmentBlocksResponse struct {
// handleGetVersion handles /version requests.
func (h *Handler) handleGetVersion(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
err := json.NewEncoder(w).Encode(struct {
Version string `json:"version"`
}{
@ -1047,6 +1121,10 @@ func errorString(err error) string {
}
func (h *Handler) handlePostClusterResizeSetCoordinator(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
// Decode request.
var req setCoordinatorRequest
err := json.NewDecoder(r.Body).Decode(&req)
@ -1084,6 +1162,10 @@ type setCoordinatorResponse struct {
// handlePostClusterResizeRemoveNode handles POST /cluster/resize/remove-node request.
func (h *Handler) handlePostClusterResizeRemoveNode(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
// Decode request.
var req removeNodeRequest
err := json.NewDecoder(r.Body).Decode(&req)
@ -1120,6 +1202,10 @@ type removeNodeResponse struct {
// handlePostClusterResizeAbort handles POST /cluster/resize/abort request.
func (h *Handler) handlePostClusterResizeAbort(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
err := h.API.ResizeAbort()
var msg string
if err != nil {
@ -1157,6 +1243,10 @@ func (h *Handler) handleRecalculateCaches(w http.ResponseWriter, r *http.Request
}
func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Request) {
if checkHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
// Verify that request is only communicating over protobufs.
if r.Header.Get("Content-Type") != "application/x-protobuf" {
http.Error(w, "Unsupported media type", http.StatusUnsupportedMediaType)

View file

@ -750,11 +750,17 @@ func TestHandler_Index_AttrStore_Diff(t *testing.T) {
blks[1].Checksum = []byte("MISMATCHED_CHECKSUM")
// Send block checksums to determine diff.
resp, err := gohttp.Post(
req, err := gohttp.NewRequest(
"POST",
s.URL+"/index/i/attr/diff",
"application/json",
strings.NewReader(`{"blocks":`+string(test.MustMarshalJSON(blks))+`}`),
)
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Accept", "application/json")
client := &gohttp.Client{}
resp, err := client.Do(req)
if err != nil {
t.Fatal(err)
}
@ -800,11 +806,17 @@ func TestHandler_Field_AttrStore_Diff(t *testing.T) {
blks[1].Checksum = []byte("MISMATCHED_CHECKSUM")
// Send block checksums to determine diff.
resp, err := gohttp.Post(
req, err := gohttp.NewRequest(
"POST",
s.URL+"/index/i/field/meta/attr/diff",
"application/json",
strings.NewReader(`{"blocks":`+string(test.MustMarshalJSON(blks))+`}`),
)
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Accept", "application/json")
client := &gohttp.Client{}
resp, err := client.Do(req)
if err != nil {
t.Fatal(err)
}

116
server.go
View file

@ -52,9 +52,9 @@ type Server struct {
closing chan struct{}
// Internal
Holder *Holder
holder *Holder
Cluster *Cluster
TranslateFile *TranslateFile
translateFile *TranslateFile
diagnostics *DiagnosticsCollector
executor *Executor
hosts []string
@ -62,11 +62,9 @@ type Server struct {
// External
handler Handler
Broadcaster Broadcaster
BroadcastReceiver BroadcastReceiver
systemInfo SystemInfo
gcNotifier GCNotifier
NewAttrStore func(string) AttrStore
logger Logger
ln net.Listener
@ -109,8 +107,7 @@ func OptServerDataDir(dir string) ServerOption {
func OptServerAttrStoreFunc(af func(string) AttrStore) ServerOption {
return func(s *Server) error {
s.NewAttrStore = af
s.Holder.NewAttrStore = af
s.holder.NewAttrStore = af
return nil
}
}
@ -182,7 +179,7 @@ func OptServerPrimaryTranslateStore(store TranslateStore) ServerOption {
func OptServerStatsClient(sc StatsClient) ServerOption {
return func(s *Server) error {
s.Holder.Stats = sc
s.holder.Stats = sc
return nil
}
}
@ -224,16 +221,13 @@ func NewServer(opts ...ServerOption) (*Server, error) {
s := &Server{
closing: make(chan struct{}),
Cluster: NewCluster(),
Holder: NewHolder(),
Broadcaster: NopBroadcaster,
holder: NewHolder(),
BroadcastReceiver: NopBroadcastReceiver,
diagnostics: NewDiagnosticsCollector(DefaultDiagnosticServer),
systemInfo: NewNopSystemInfo(),
gcNotifier: NopGCNotifier,
NewAttrStore: NewNopAttrStore,
antiEntropyInterval: time.Minute * 10,
metricInterval: 0,
diagnosticInterval: 0,
@ -254,19 +248,19 @@ func NewServer(opts ...ServerOption) (*Server, error) {
return nil, err
}
s.Holder.Path = path
s.Holder.Logger = s.logger
s.Holder.Stats.SetLogger(s.logger)
s.holder.Path = path
s.holder.Logger = s.logger
s.holder.Stats.SetLogger(s.logger)
s.Cluster.Path = path
s.Cluster.Logger = s.logger
s.Cluster.Holder = s.Holder
s.Cluster.Holder = s.holder
// Initialize translation database.
s.TranslateFile = NewTranslateFile()
s.TranslateFile.Path = filepath.Join(path, "keys")
s.TranslateFile.PrimaryTranslateStore = s.primaryTranslateStore
if err := s.TranslateFile.Open(); err != nil {
s.translateFile = NewTranslateFile()
s.translateFile.Path = filepath.Join(path, "keys")
s.translateFile.PrimaryTranslateStore = s.primaryTranslateStore
if err := s.translateFile.Open(); err != nil {
return nil, err
}
@ -292,15 +286,15 @@ func NewServer(opts ...ServerOption) (*Server, error) {
}
// Append the NodeID tag to stats.
s.Holder.Stats = s.Holder.Stats.WithTags(fmt.Sprintf("NodeID:%s", s.NodeID))
s.holder.Stats = s.holder.Stats.WithTags(fmt.Sprintf("NodeID:%s", s.NodeID))
s.executor.Holder = s.Holder
s.executor.Holder = s.holder
s.executor.Node = node
s.executor.Cluster = s.Cluster
s.executor.TranslateStore = s.TranslateFile
s.executor.TranslateStore = s.translateFile
s.executor.MaxWritesPerRequest = s.maxWritesPerRequest
s.handler.GetAPI().Executor = s.executor
s.handler.GetAPI().TranslateStore = s.TranslateFile
s.handler.GetAPI().TranslateStore = s.translateFile
return s, nil
}
@ -313,25 +307,25 @@ func (s *Server) Open() error {
}
// Log startup
err := s.Holder.logStartup()
err := s.holder.logStartup()
if err != nil {
log.Println(errors.Wrap(err, "logging startup"))
}
// Cluster settings.
s.Cluster.Broadcaster = s.Broadcaster
s.Cluster.Broadcaster = s
s.Cluster.MaxWritesPerRequest = s.maxWritesPerRequest
// Initialize HTTP handler.
api := s.handler.GetAPI()
api.Holder = s.Holder
api.Broadcaster = s.Broadcaster
api.Holder = s.holder
api.Broadcaster = s
api.BroadcastHandler = s
api.StatusHandler = s
api.Cluster = s.Cluster
// Initialize Holder.
s.Holder.Broadcaster = s.Broadcaster
s.holder.Broadcaster = s
// Serve handler.
go s.handler.Serve(s.ln, s.closing)
@ -347,7 +341,7 @@ func (s *Server) Open() error {
}
// Open holder.
if err := s.Holder.Open(); err != nil {
if err := s.holder.Open(); err != nil {
return fmt.Errorf("opening Holder: %v", err)
}
if err := s.Cluster.setNodeState(NodeStateReady); err != nil {
@ -382,11 +376,11 @@ func (s *Server) Close() error {
if s.Cluster != nil {
s.Cluster.close()
}
if s.Holder != nil {
s.Holder.Close()
if s.holder != nil {
s.holder.Close()
}
if s.TranslateFile != nil {
s.TranslateFile.Close()
if s.translateFile != nil {
s.translateFile.Close()
}
return nil
@ -398,7 +392,7 @@ func (s *Server) LoadNodeID() string {
if s.NodeID != "" {
return s.NodeID
}
nodeID, err := s.Holder.loadNodeID()
nodeID, err := s.holder.loadNodeID()
if err != nil {
s.logger.Printf("loading NodeID: %v", err)
return s.NodeID
@ -426,18 +420,18 @@ func (s *Server) monitorAntiEntropy() {
case <-s.closing:
return
case <-ticker.C:
s.Holder.Stats.Count("AntiEntropy", 1, 1.0)
s.holder.Stats.Count("AntiEntropy", 1, 1.0)
}
t := time.Now()
s.logger.Printf("holder sync beginning")
// Initialize syncer with local holder and remote client.
var syncer HolderSyncer
syncer.Holder = s.Holder
syncer.Holder = s.holder
syncer.Node = s.Cluster.Node
syncer.Cluster = s.Cluster
syncer.Closing = s.closing
syncer.Stats = s.Holder.Stats.WithTags("HolderSyncer")
syncer.Stats = s.holder.Stats.WithTags("HolderSyncer")
// Sync holders.
if err := syncer.SyncHolder(); err != nil {
@ -448,7 +442,7 @@ func (s *Server) monitorAntiEntropy() {
// Record successful sync in log.
s.logger.Printf("holder sync complete")
dif := time.Since(t)
s.Holder.Stats.Histogram("AntiEntropyDuration", float64(dif), 1.0)
s.holder.Stats.Histogram("AntiEntropyDuration", float64(dif), 1.0)
}
}
@ -456,23 +450,23 @@ func (s *Server) monitorAntiEntropy() {
func (s *Server) ReceiveMessage(pb proto.Message) error {
switch obj := pb.(type) {
case *internal.CreateSliceMessage:
idx := s.Holder.Index(obj.Index)
idx := s.holder.Index(obj.Index)
if idx == nil {
return fmt.Errorf("Local Index not found: %s", obj.Index)
}
idx.SetRemoteMaxSlice(obj.Slice)
case *internal.CreateIndexMessage:
opt := IndexOptions{}
_, err := s.Holder.CreateIndex(obj.Index, opt)
_, err := s.holder.CreateIndex(obj.Index, opt)
if err != nil {
return err
}
case *internal.DeleteIndexMessage:
if err := s.Holder.DeleteIndex(obj.Index); err != nil {
if err := s.holder.DeleteIndex(obj.Index); err != nil {
return err
}
case *internal.CreateFieldMessage:
idx := s.Holder.Index(obj.Index)
idx := s.holder.Index(obj.Index)
if idx == nil {
return fmt.Errorf("Local Index not found: %s", obj.Index)
}
@ -482,12 +476,12 @@ func (s *Server) ReceiveMessage(pb proto.Message) error {
return err
}
case *internal.DeleteFieldMessage:
idx := s.Holder.Index(obj.Index)
idx := s.holder.Index(obj.Index)
if err := idx.DeleteField(obj.Field); err != nil {
return err
}
case *internal.CreateViewMessage:
f := s.Holder.Field(obj.Index, obj.Field)
f := s.holder.Field(obj.Index, obj.Field)
if f == nil {
return fmt.Errorf("Local Field not found: %s", obj.Field)
}
@ -496,7 +490,7 @@ func (s *Server) ReceiveMessage(pb proto.Message) error {
return err
}
case *internal.DeleteViewMessage:
f := s.Holder.Field(obj.Index, obj.Field)
f := s.holder.Field(obj.Index, obj.Field)
if f == nil {
return fmt.Errorf("Local Field not found: %s", obj.Field)
}
@ -529,7 +523,7 @@ func (s *Server) ReceiveMessage(pb proto.Message) error {
return err
}
case *internal.RecalculateCaches:
s.Holder.RecalculateCaches()
s.holder.RecalculateCaches()
case *internal.NodeEventMessage:
s.Cluster.ReceiveEvent(DecodeNodeEvent(obj))
}
@ -581,14 +575,14 @@ func (s *Server) LocalStatus() (proto.Message, error) {
if s.Cluster == nil {
return nil, errors.New("Server.Cluster is nil")
}
if s.Holder == nil {
if s.holder == nil {
return nil, errors.New("Server.Holder is nil")
}
ns := internal.NodeStatus{
Node: EncodeNode(s.Cluster.Node),
MaxSlices: s.Holder.EncodeMaxSlices(),
Schema: s.Holder.EncodeSchema(),
MaxSlices: s.holder.EncodeMaxSlices(),
Schema: s.holder.EncodeSchema(),
}
return &ns, nil
@ -608,7 +602,7 @@ func (s *Server) HandleRemoteStatus(pb proto.Message) error {
go func() {
// Make sure the holder has opened.
<-s.Holder.opened
<-s.holder.opened
err := s.mergeRemoteStatus(pb.(*internal.NodeStatus))
if err != nil {
@ -626,14 +620,14 @@ func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error {
}
// Sync schema.
if err := s.Holder.ApplySchema(ns.Schema); err != nil {
if err := s.holder.ApplySchema(ns.Schema); err != nil {
return errors.Wrap(err, "applying schema")
}
// Sync maxSlices.
oldmaxslices := s.Holder.MaxSlices()
oldmaxslices := s.holder.MaxSlices()
for index, newMax := range ns.MaxSlices.Standard {
localIndex := s.Holder.Index(index)
localIndex := s.holder.Index(index)
// if we don't know about an index locally, log an error because
// indexes should be created and synced prior to slice creation
if localIndex == nil {
@ -721,26 +715,26 @@ func (s *Server) monitorRuntime() {
return
case <-s.gcNotifier.AfterGC():
// GC just ran.
s.Holder.Stats.Count("garbage_collection", 1, 1.0)
s.holder.Stats.Count("garbage_collection", 1, 1.0)
case <-ticker.C:
}
// Record the number of go routines.
s.Holder.Stats.Gauge("goroutines", float64(runtime.NumGoroutine()), 1.0)
s.holder.Stats.Gauge("goroutines", float64(runtime.NumGoroutine()), 1.0)
openFiles, err := countOpenFiles()
// Open File handles.
if err == nil {
s.Holder.Stats.Gauge("OpenFiles", float64(openFiles), 1.0)
s.holder.Stats.Gauge("OpenFiles", float64(openFiles), 1.0)
}
// Runtime memory metrics.
runtime.ReadMemStats(&m)
s.Holder.Stats.Gauge("HeapAlloc", float64(m.HeapAlloc), 1.0)
s.Holder.Stats.Gauge("HeapInuse", float64(m.HeapInuse), 1.0)
s.Holder.Stats.Gauge("StackInuse", float64(m.StackInuse), 1.0)
s.Holder.Stats.Gauge("Mallocs", float64(m.Mallocs), 1.0)
s.Holder.Stats.Gauge("Frees", float64(m.Frees), 1.0)
s.holder.Stats.Gauge("HeapAlloc", float64(m.HeapAlloc), 1.0)
s.holder.Stats.Gauge("HeapInuse", float64(m.HeapInuse), 1.0)
s.holder.Stats.Gauge("StackInuse", float64(m.StackInuse), 1.0)
s.holder.Stats.Gauge("Mallocs", float64(m.Mallocs), 1.0)
s.holder.Stats.Gauge("Frees", float64(m.Frees), 1.0)
}
}

View file

@ -65,10 +65,10 @@ type Command struct {
// Standard input/output
*pilosa.CmdIO
// Started will be closed once Command.Run is finished.
// Started will be closed once Command.Start is finished.
Started chan struct{}
// Done will be closed when Command.Close() is called
Done chan struct{}
// done will be closed when Command.Close() is called
done chan struct{}
// Passed to the Gossip implementation.
logOutput io.Writer
@ -83,7 +83,7 @@ func NewCommand(stdin io.Reader, stdout, stderr io.Writer) *Command {
CmdIO: pilosa.NewCmdIO(stdin, stdout, stderr),
Started: make(chan struct{}),
Done: make(chan struct{}),
done: make(chan struct{}),
}
}
@ -125,7 +125,7 @@ func (m *Command) Wait() error {
// Second signal causes a hard shutdown.
go func() { <-c; os.Exit(1) }()
return errors.Wrap(m.Close(), "closing command")
case <-m.Done:
case <-m.done:
m.logger.Printf("Server closed externally")
return nil
}
@ -293,7 +293,6 @@ func (m *Command) SetupNetworking() error {
}
gossipMemberSet.Logger = m.logger
m.Server.Cluster.MemberSet = gossipMemberSet
m.Server.Broadcaster = m.Server
m.Server.BroadcastReceiver = gossipMemberSet
return nil
}
@ -305,7 +304,7 @@ func (m *Command) Close() error {
if closer, ok := m.logOutput.(io.Closer); ok {
logErr = closer.Close()
}
close(m.Done)
close(m.done)
if serveErr != nil && logErr != nil {
return fmt.Errorf("closing server: '%v', closing logs: '%v'", serveErr, logErr)
} else if logErr != nil {

View file

@ -148,6 +148,7 @@ func MustParseURLHost(rawurl string) string {
// MustNewHTTPRequest creates a new HTTP request. Panic on error.
func MustNewHTTPRequest(method, urlStr string, body io.Reader) *gohttp.Request {
req, err := gohttp.NewRequest(method, urlStr, body)
req.Header.Add("Accept", "application/json")
if err != nil {
panic(err)
}

View file

@ -25,7 +25,6 @@ import (
"testing"
"time"
"github.com/pilosa/pilosa/boltdb"
"github.com/pilosa/pilosa/gossip"
"github.com/pilosa/pilosa/http"
"github.com/pilosa/pilosa/server"
@ -164,9 +163,6 @@ func (m *Main) Reopen() error {
return errors.Wrap(err, "setting up server")
}
m.Server.NewAttrStore = boltdb.NewAttrStore
m.Server.Holder.NewAttrStore = m.Server.NewAttrStore
// Run new program.
if err := m.Start(); err != nil {
return err
@ -267,10 +263,15 @@ func (m *Main) RecalculateCaches() error {
// MustDo executes http.Do() with an http.NewRequest(). Panic on error.
func MustDo(method, urlStr string, body string) *httpResponse {
req, err := gohttp.NewRequest(method, urlStr, strings.NewReader(body))
if err != nil {
panic(err)
}
req, err := gohttp.NewRequest(
method,
urlStr,
strings.NewReader(body),
)
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Accept", "application/json")
resp, err := gohttp.DefaultClient.Do(req)
if err != nil {
panic(err)

View file

@ -17,6 +17,7 @@ package test_test
import (
"encoding/json"
"net/http"
"strings"
"testing"
"github.com/pilosa/pilosa"
@ -32,12 +33,21 @@ func TestNewCluster(t *testing.T) {
t.Fatalf("node %d does not have the same coordinator as node 0. '%v' and '%v' respectively", i, coordi, coordinator)
}
}
req, err := http.NewRequest(
"GET",
"http://"+cluster[0].Server.Addr().String()+"/status",
strings.NewReader(""),
)
response, err := http.Get("http://" + cluster[0].Server.Addr().String() + "/status")
req.Header.Set("Accept", "application/json")
resp, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatalf("getting schema: %v", err)
t.Fatalf("sending request: %v", err)
}
dec := json.NewDecoder(response.Body)
defer resp.Body.Close()
dec := json.NewDecoder(resp.Body)
body := struct {
State string
Nodes []struct {