diff --git a/api.go b/api.go new file mode 100644 index 000000000..82ab3a4ce --- /dev/null +++ b/api.go @@ -0,0 +1,896 @@ +// Copyright 2017 Pilosa Corp. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package pilosa + +import ( + "context" + "encoding/csv" + "fmt" + "io" + "net/http" + "reflect" + "strconv" + "strings" + "time" + + "github.com/gogo/protobuf/proto" + "github.com/pilosa/pilosa/internal" + "github.com/pilosa/pilosa/pql" + "github.com/pkg/errors" +) + +type API struct { + Holder *Holder + // The execution engine for running queries. + Executor interface { + Execute(context context.Context, index string, query *pql.Query, slices []uint64, opt *ExecOptions) ([]interface{}, error) + } + Broadcaster Broadcaster + BroadcastHandler BroadcastHandler + StatusHandler StatusHandler + Cluster *Cluster + URI URI + RemoteClient *http.Client + Logger Logger +} + +func NewAPI() *API { + return &API{ + Broadcaster: NopBroadcaster, + //BroadcastHandler: NopBroadcastHandler, // TODO: implement the nop + //StatusHandler: NopStatusHandler, // TODO: implement the nop + Logger: NopLogger, + } +} + +func (a *API) ExecuteQuery(ctx context.Context, req *QueryRequest) (QueryResponse, error) { + resp := QueryResponse{} + + q, err := pql.NewParser(strings.NewReader(req.Query)).Parse() + if err != nil { + return resp, err + } + execOpts := &ExecOptions{ + Remote: req.Remote, + ExcludeAttrs: req.ExcludeAttrs, + ExcludeBits: req.ExcludeBits, + } + results, err := a.Executor.Execute(ctx, req.Index, q, req.Slices, execOpts) + if err != nil { + return resp, err + } + resp.Results = results + + // Fill column attributes if requested. + if req.ColumnAttrs && !req.ExcludeBits { + // Consolidate all column ids across all calls. + var columnIDs []uint64 + for _, result := range results { + bm, ok := result.(*Bitmap) + if !ok { + continue + } + columnIDs = uint64Slice(columnIDs).merge(bm.Bits()) + } + + // Retrieve column attributes across all calls. + columnAttrSets, err := a.readColumnAttrSets(a.Holder.Index(req.Index), columnIDs) + if err != nil { + return resp, err + } + resp.ColumnAttrSets = columnAttrSets + } + return resp, nil +} + +// readColumnAttrSets returns a list of column attribute objects by id. +func (api *API) readColumnAttrSets(index *Index, ids []uint64) ([]*ColumnAttrSet, error) { + if index == nil { + return nil, nil + } + + ax := make([]*ColumnAttrSet, 0, len(ids)) + for _, id := range ids { + // Read attributes for column. Skip column if empty. + attrs, err := index.ColumnAttrStore().Attrs(id) + if err != nil { + return nil, err + } else if len(attrs) == 0 { + continue + } + + // Append column with attributes. + ax = append(ax, &ColumnAttrSet{ID: id, Attrs: attrs}) + } + + return ax, nil +} + +func (api *API) CreateIndex(ctx context.Context, indexName string, options IndexOptions) (*Index, error) { + // Create index. + index, err := api.Holder.CreateIndex(indexName, options) + if err != nil { + return nil, err + } + // Send the create index message to all nodes. + err = api.Broadcaster.SendSync( + &internal.CreateIndexMessage{ + Index: indexName, + Meta: options.Encode(), + }) + if err != nil { + api.Logger.Printf("problem sending CreateIndex message: %s", err) + return nil, err + } + api.Holder.Stats.Count("createIndex", 1, 1.0) + return index, nil +} + +func (api *API) ReadIndex(ctx context.Context, indexName string) (*Index, error) { + index := api.Holder.Index(indexName) + if index == nil { + return nil, ErrIndexNotFound + } + return index, nil +} + +func (api *API) DeleteIndex(ctx context.Context, indexName string) error { + // Delete index from the holder. + err := api.Holder.DeleteIndex(indexName) + if err != nil { + return err + } + // Send the delete index message to all nodes. + err = api.Broadcaster.SendSync( + &internal.DeleteIndexMessage{ + Index: indexName, + }) + if err != nil { + api.Logger.Printf("problem sending DeleteIndex message: %s", err) + return err + } + api.Holder.Stats.Count("deleteIndex", 1, 1.0) + return nil +} + +func (api *API) CreateFrame(ctx context.Context, indexName string, frameName string, options FrameOptions) (*Frame, error) { + // Find index. + index := api.Holder.Index(indexName) + if index == nil { + return nil, ErrIndexNotFound + } + + // Create frame. + frame, err := index.CreateFrame(frameName, options) + if err != nil { + return nil, err + } + + // Send the create frame message to all nodes. + err = api.Broadcaster.SendSync( + &internal.CreateFrameMessage{ + Index: indexName, + Frame: frameName, + Meta: options.Encode(), + }) + if err != nil { + api.Logger.Printf("problem sending CreateFrame message: %s", err) + return nil, err + } + api.Holder.Stats.CountWithCustomTags("createFrame", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) + return frame, nil +} + +func (api *API) DeleteFrame(ctx context.Context, indexName string, frameName string) error { + // Find index. + index := api.Holder.Index(indexName) + if index == nil { + return ErrIndexNotFound + } + + // Delete frame from the index. + if err := index.DeleteFrame(frameName); err != nil { + return err + } + + // Send the delete frame message to all nodes. + err := api.Broadcaster.SendSync( + &internal.DeleteFrameMessage{ + Index: indexName, + Frame: frameName, + }) + if err != nil { + api.Logger.Printf("problem sending DeleteFrame message: %s", err) + return err + } + api.Holder.Stats.CountWithCustomTags("deleteFrame", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) + return nil +} + +func (api *API) ExportCSV(ctx context.Context, indexName string, frameName string, viewName string, slice uint64, w io.Writer) error { + // Validate that this handler owns the slice. + if !api.Cluster.OwnsFragment(api.LocalID(), indexName, slice) { + api.Logger.Printf("host does not own slice %s-%s slice:%d", api.URI, indexName, slice) + return ErrClusterDoesNotOwnSlice + } + + // Find the fragment. + f := api.Holder.Fragment(indexName, frameName, viewName, slice) + if f == nil { + return ErrFragmentNotFound + } + + // Wrap writer with a CSV writer. + cw := csv.NewWriter(w) + + // Iterate over each bit. + if err := f.ForEachBit(func(rowID, columnID uint64) error { + return cw.Write([]string{ + strconv.FormatUint(rowID, 10), + strconv.FormatUint(columnID, 10), + }) + }); err != nil { + return err + } + + // Ensure data is flushed. + cw.Flush() + + return nil +} + +func (api *API) FragmentNodes(ctx context.Context, indexName string, slice uint64) []*Node { + return api.Cluster.FragmentNodes(indexName, slice) +} + +func (api *API) FragmentData(ctx context.Context, indexName string, frameName string, viewName string, slice uint64) (*Fragment, error) { + // Retrieve fragment from holder. + f := api.Holder.Fragment(indexName, frameName, viewName, slice) + if f == nil { + return nil, ErrFragmentNotFound + } + return f, nil +} + +func (api *API) WriteFragmentData(ctx context.Context, indexName string, frameName string, viewName string, slice uint64, reader io.ReadCloser) error { + // Retrieve frame. + f := api.Holder.Frame(indexName, frameName) + if f == nil { + return ErrFrameNotFound + } + + // Retrieve view. + view, err := f.CreateViewIfNotExists(viewName) + if err != nil { + return err + } + + // Retrieve fragment from frame. + frag, err := view.CreateFragmentIfNotExists(slice) + if err != nil { + return err + } + + // Read fragment in from request body. + if _, err := frag.ReadFrom(reader); err != nil { + return err + } + return nil +} + +func (api *API) FragmentBlockData(ctx context.Context, req internal.BlockDataRequest) (internal.BlockDataResponse, error) { + // Retrieve fragment from holder. + f := api.Holder.Fragment(req.Index, req.Frame, req.View, req.Slice) + if f == nil { + return internal.BlockDataResponse{}, ErrFragmentNotFound + } + + // Read data + var resp internal.BlockDataResponse + resp.RowIDs, resp.ColumnIDs = f.BlockData(int(req.Block)) + return resp, nil +} + +func (api *API) FragmentBlocks(ctx context.Context, indexName string, frameName string, viewName string, slice uint64) ([]FragmentBlock, error) { + // Retrieve fragment from holder. + f := api.Holder.Fragment(indexName, frameName, viewName, slice) + if f == nil { + return nil, ErrFragmentNotFound + } + + // Retrieve blocks. + blocks := f.Blocks() + return blocks, nil +} + +func (api *API) RestoreFrame(ctx context.Context, indexName string, frameName string, host *URI) error { + // Create a client for the remote cluster. + client := NewInternalHTTPClientFromURI(host, api.RemoteClient) + + // Determine the maximum number of slices. + maxSlices, err := client.MaxSliceByIndex(ctx) + if err != nil { + return err + } + + // Retrieve frame. + f := api.Holder.Frame(indexName, frameName) + if f == nil { + return ErrFrameNotFound + } + + // Retrieve list of all views. + views, err := client.FrameViews(ctx, indexName, frameName) + if err != nil { + return err + } + + // Loop over each slice and import it if this node owns it. + for slice := uint64(0); slice <= maxSlices[indexName]; slice++ { + // Ignore this slice if we don't own it. + if !api.Cluster.OwnsFragment(api.LocalID(), indexName, slice) { + continue + } + + // Loop over view names. + for _, view := range views { + // Create view. + v, err := f.CreateViewIfNotExists(view) + if err != nil { + return err + } + + // Otherwise retrieve the local fragment. + frag, err := v.CreateFragmentIfNotExists(slice) + if err != nil { + return err + } + + // Stream backup from remote node. + rd, err := client.BackupSlice(ctx, indexName, frameName, view, slice) + if err != nil { + return err + } else if rd == nil { + continue // slice doesn't exist + } + + // Restore to local frame and always close reader. + if err := func() error { + defer rd.Close() + if _, err := frag.ReadFrom(rd); err != nil { + return err + } + return nil + }(); err != nil { + return err + } + } + } + + return nil +} + +func (api *API) ClusterHosts(ctx context.Context) []*Node { + return api.Cluster.Nodes +} + +func (api *API) CreateInputDefinition(ctx context.Context, indexName string, inputDefName string, inputDef InputDefinitionInfo) error { + // Find index. + index := api.Holder.Index(indexName) + if index == nil { + return ErrIndexNotFound + } + + if err := inputDef.Validate(); err != nil { + return err + } + + // Encode InputDefinition to its internal representation. + def := inputDef.Encode() + def.Name = inputDefName + + // Create InputDefinition. + if _, err := index.CreateInputDefinition(def); err != nil { + return err + } + + err := api.Broadcaster.SendSync( + &internal.CreateInputDefinitionMessage{ + Index: indexName, + Definition: def, + }) + if err != nil { + api.Logger.Printf("problem sending CreateInputDefinition message: %s", err) + } + return nil +} + +func (api *API) InputDefinition(ctx context.Context, indexName string, inputDefName string) (*InputDefinition, error) { + // Find index. + index := api.Holder.Index(indexName) + if index == nil { + return nil, ErrIndexNotFound + } + + inputDef, err := index.InputDefinition(inputDefName) + if err != nil { + return nil, err + } + return inputDef, nil +} + +func (api *API) DeleteInputDefinition(ctx context.Context, indexName string, inputDefName string) error { + // Find index. + index := api.Holder.Index(indexName) + if index == nil { + return ErrIndexNotFound + } + + // Delete input definition from the index. + if err := index.DeleteInputDefinition(inputDefName); err != nil { + return err + } + + err := api.Broadcaster.SendSync( + &internal.DeleteInputDefinitionMessage{ + Index: indexName, + Name: inputDefName, + }) + if err != nil { + api.Logger.Printf("problem sending DeleteInputDefinition message: %s", err) + } + return nil +} + +func (api *API) WriteInput(ctx context.Context, indexName string, inputDefName string, reqs []interface{}) error { + // Find index. + index := api.Holder.Index(indexName) + if index == nil { + return ErrIndexNotFound + } + + for _, req := range reqs { + bits, err := api.inputJSONDataParser(req.(map[string]interface{}), index, inputDefName) + if err != nil { + return err + } + for fr, bs := range bits { + if err := index.InputBits(fr, bs); err != nil { + return err + } + } + } + + return nil +} + +func (api *API) RecalculateCaches(ctx context.Context) error { + err := api.Broadcaster.SendSync(&internal.RecalculateCaches{}) + if err != nil { + return errors.Wrap(err, "broacasting message") + } + api.Holder.RecalculateCaches() + return nil +} + +func (api *API) PostClusterMessage(ctx context.Context, pb proto.Message) error { + // Forward the error message. + if err := api.BroadcastHandler.ReceiveMessage(pb); err != nil { + return err + } + return nil +} + +func (api *API) LocalID() string { + return api.Cluster.Node.ID +} + +func (api *API) Schema(ctx context.Context) []*IndexInfo { + return api.Holder.Schema() +} + +func (api *API) Status(ctx context.Context) (proto.Message, error) { + return api.StatusHandler.ClusterStatus() +} + +func (api *API) CreateFrameField(ctx context.Context, indexName string, frameName string, field *Field) error { + // Retrieve frame by name. + f := api.Holder.Frame(indexName, frameName) + if f == nil { + return ErrFrameNotFound + } + + // Create new field. + if err := f.CreateField(field); err != nil { + return err + } + + // Send the create field message to all nodes. + err := api.Broadcaster.SendSync( + &internal.CreateFieldMessage{ + Index: indexName, + Frame: frameName, + Field: encodeField(field), + }) + if err != nil { + api.Logger.Printf("problem sending CreateField message: %s", err) + } + return err +} + +func (api *API) DeleteFrameField(ctx context.Context, indexName string, frameName string, fieldName string) error { + // Retrieve frame by name. + f := api.Holder.Frame(indexName, frameName) + if f == nil { + return ErrFrameNotFound + } + + // Delete field. + if err := f.DeleteField(fieldName); err != nil { + return err + } + + // Send the delete field message to all nodes. + err := api.Broadcaster.SendSync( + &internal.DeleteFieldMessage{ + Index: indexName, + Frame: frameName, + Field: fieldName, + }) + if err != nil { + api.Logger.Printf("problem sending DeleteField message: %s", err) + } + return err +} + +func (api *API) FrameFields(ctx context.Context, indexName string, frameName string) ([]*Field, error) { + index := api.Holder.index(indexName) + if index == nil { + return nil, ErrIndexNotFound + } + + frame := index.frame(frameName) + if frame == nil { + return nil, ErrFrameNotFound + } + + return frame.GetFields() +} + +func (api *API) FrameViews(ctx context.Context, indexName string, frameName string) ([]*View, error) { + // Retrieve views. + f := api.Holder.Frame(indexName, frameName) + if f == nil { + return nil, ErrFrameNotFound + } + + // Fetch views. + views := f.Views() + return views, nil +} + +func (api *API) DeleteView(ctx context.Context, indexName string, frameName string, viewName string) error { + // Retrieve frame. + f := api.Holder.Frame(indexName, frameName) + if f == nil { + return ErrFrameNotFound + } + + // Delete the view. + if err := f.DeleteView(viewName); err != nil { + // Ingore this error becuase views do not exist on all nodes due to slice distribution. + if err != ErrInvalidView { + return err + } + } + + // Send the delete view message to all nodes. + err := api.Broadcaster.SendSync( + &internal.DeleteViewMessage{ + Index: indexName, + Frame: frameName, + View: viewName, + }) + if err != nil { + api.Logger.Printf("problem sending DeleteView message: %s", err) + } + + return err +} + +func (api *API) IndexAttrDiff(ctx context.Context, indexName string, blocks []AttrBlock) (map[uint64]map[string]interface{}, error) { + // Retrieve index from holder. + index := api.Holder.Index(indexName) + if index == nil { + return nil, ErrIndexNotFound + } + + // Retrieve local blocks. + localBlocks, err := index.ColumnAttrStore().Blocks() + if err != nil { + return nil, err + } + + // Read all attributes from all mismatched blocks. + attrs := make(map[uint64]map[string]interface{}) + for _, blockID := range AttrBlocks(localBlocks).Diff(blocks) { + // Retrieve block data. + m, err := index.ColumnAttrStore().BlockData(blockID) + if err != nil { + return nil, err + } + + // Copy to index-wide struct. + for k, v := range m { + attrs[k] = v + } + } + return attrs, nil +} + +func (api *API) FrameAttrDiff(ctx context.Context, indexName string, frameName string, blocks []AttrBlock) (map[uint64]map[string]interface{}, error) { + // Retrieve index from holder. + f := api.Holder.Frame(indexName, frameName) + if f == nil { + return nil, ErrFrameNotFound + } + + // Retrieve local blocks. + localBlocks, err := f.RowAttrStore().Blocks() + if err != nil { + return nil, err + } + + // Read all attributes from all mismatched blocks. + attrs := make(map[uint64]map[string]interface{}) + for _, blockID := range AttrBlocks(localBlocks).Diff(blocks) { + // Retrieve block data. + m, err := f.RowAttrStore().BlockData(blockID) + if err != nil { + return nil, err + } + + // Copy to index-wide struct. + for k, v := range m { + attrs[k] = v + } + } + return attrs, nil +} + +func (api *API) Import(ctx context.Context, req internal.ImportRequest) error { + _, frame, err := api.indexFrame(req.Index, req.Frame, req.Slice) + if err != nil { + return err + } + + // Convert timestamps to time.Time. + timestamps := make([]*time.Time, len(req.Timestamps)) + for i, ts := range req.Timestamps { + if ts == 0 { + continue + } + t := time.Unix(0, ts) + timestamps[i] = &t + } + + // Import into fragment. + err = frame.Import(req.RowIDs, req.ColumnIDs, timestamps) + if err != nil { + api.Logger.Printf("import error: index=%s, frame=%s, slice=%d, bits=%d, err=%s", req.Index, req.Frame, req.Slice, len(req.ColumnIDs), err) + } + return err +} + +func (api *API) ImportValue(ctx context.Context, req internal.ImportValueRequest) error { + _, frame, err := api.indexFrame(req.Index, req.Frame, req.Slice) + if err != nil { + return err + } + + // Import into fragment. + err = frame.ImportValue(req.Field, req.ColumnIDs, req.Values) + if err != nil { + api.Logger.Printf("import error: index=%s, frame=%s, slice=%d, field=%s, bits=%d, err=%s", req.Index, req.Frame, req.Slice, req.Field, len(req.ColumnIDs), err) + } + return err +} + +func (api *API) ModifyIndexTimeQuantum(ctx context.Context, indexName string, timeQuantum TimeQuantum) error { + // Retrieve index by name. + index := api.Holder.Index(indexName) + if index == nil { + return ErrIndexNotFound + } + + // Set default time quantum on index. + return index.SetTimeQuantum(timeQuantum) +} + +func (api *API) ModifyFrameTimeQuantum(ctx context.Context, indexName string, frameName string, timeQuantum TimeQuantum) error { + // Retrieve index by name. + frame := api.Holder.Frame(indexName, frameName) + if frame == nil { + return ErrFrameNotFound + } + + // Set default time quantum on index. + return frame.SetTimeQuantum(timeQuantum) +} + +func (api *API) MaxSlices(ctx context.Context) map[string]uint64 { + return api.Holder.MaxSlices() +} + +func (api *API) MaxInverseSlices(ctx context.Context) map[string]uint64 { + return api.Holder.MaxInverseSlices() +} + +func (api *API) StatsWithTags(tags []string) StatsClient { + if api.Holder == nil || api.Cluster == nil { + return nil + } + return api.Holder.Stats.WithTags(tags...) +} + +func (api *API) ClusterLongQueryTime() time.Duration { + if api.Cluster == nil { + return 0 + } + return api.Cluster.LongQueryTime +} + +func (api *API) indexFrame(indexName string, frameName string, slice uint64) (*Index, *Frame, error) { + // Validate that this handler owns the slice. + if !api.Cluster.OwnsFragment(api.LocalID(), indexName, slice) { + api.Logger.Printf("host does not own slice %s-%s slice:%d", api.URI, indexName, slice) + return nil, nil, ErrClusterDoesNotOwnSlice + } + + // Find the Index. + api.Logger.Printf("importing: %v %v %v", indexName, frameName, slice) + index := api.Holder.Index(indexName) + if index == nil { + api.Logger.Printf("fragment error: index=%s, frame=%s, slice=%d, err=%s", indexName, frameName, slice, ErrIndexNotFound.Error()) + return nil, nil, ErrIndexNotFound + } + + // Retrieve frame. + frame := index.Frame(frameName) + if frame == nil { + api.Logger.Printf("frame error: index=%s, frame=%s, slice=%d, err=%s", indexName, frameName, slice, ErrFrameNotFound.Error()) + return nil, nil, ErrFrameNotFound + } + return index, frame, nil +} + +// inputJSONDataParser validates input json file and executes SetBit. +func (api *API) inputJSONDataParser(req map[string]interface{}, index *Index, name string) (map[string][]*Bit, error) { + inputDef, err := index.InputDefinition(name) + if err != nil { + return nil, err + } + // If field in input data is not in defined definition, return error. + var colValue uint64 + validFields := make(map[string]bool) + timestampFrame := make(map[string]int64) + for _, field := range inputDef.Fields() { + validFields[field.Name] = true + if field.PrimaryKey { + value, ok := req[field.Name] + if !ok { + return nil, fmt.Errorf("primary key does not exist") + } + rawValue, ok := value.(float64) // The default JSON marshalling will interpret this as a float + if !ok { + return nil, fmt.Errorf("float64 require, got value:%s, type: %s", value, reflect.TypeOf(value)) + } + colValue = uint64(rawValue) + } + // Find frame that need to add timestamp. + for _, action := range field.Actions { + if action.ValueDestination == InputSetTimestamp { + timestampFrame[action.Frame], err = GetTimeStamp(req, field.Name) + if err != nil { + return nil, err + } + } + } + } + + for key := range req { + _, ok := validFields[key] + if !ok { + return nil, fmt.Errorf("field not found: %s", key) + } + } + + setBits := make(map[string][]*Bit) + + for _, field := range inputDef.Fields() { + // skip field that defined in definition but not in input data + if _, ok := req[field.Name]; !ok { + continue + } + + // Looking into timestampFrame map and set timestamp to the whole frame + for _, action := range field.Actions { + frame := action.Frame + timestamp := timestampFrame[action.Frame] + // Skip input data field values that are set to null + if req[field.Name] == nil { + continue + } + bit, err := HandleAction(action, req[field.Name], colValue, timestamp) + if err != nil { + return nil, fmt.Errorf("error handling action: %s, err: %s", action.ValueDestination, err) + } + if bit != nil { + setBits[frame] = append(setBits[frame], bit) + } + } + } + return setBits, nil +} + +func (api *API) SetCoordinator(ctx context.Context, id string) (oldNode, newNode *Node, err error) { + oldNode = api.Cluster.nodeByID(api.Cluster.Coordinator) + newNode = api.Cluster.nodeByID(id) + if newNode == nil { + return nil, nil, errors.Wrap(ErrNodeIDNotExists, "getting new node") + } + + // If the new coordinator is this node, do the SetCoordinator directly. + if newNode.ID == api.LocalID() { + return oldNode, newNode, api.Cluster.SetCoordinator(newNode) + } + + // Send the set-coordinator message to new node. + err = api.Broadcaster.SendTo( + newNode, + &internal.SetCoordinatorMessage{ + New: EncodeNode(newNode), + }) + if err != nil { + return nil, nil, fmt.Errorf("problem sending SetCoordinator message: %s", err) + } + return oldNode, newNode, nil +} + +func (api *API) RemoveNode(id string) (*Node, error) { + removeNode := api.Cluster.nodeByID(id) + if removeNode == nil { + return nil, errors.Wrap(ErrNodeIDNotExists, "finding node to remove") + } + + // Start the resize process (similar to NodeJoin) + err := api.Cluster.NodeLeave(removeNode) + if err != nil { + return removeNode, errors.Wrap(err, "calling node leave") + } + return removeNode, nil +} + +func (api *API) ResizeAbort() error { + if !api.Cluster.IsCoordinator() { + return ErrNodeNotCoordinator + } + err := api.Cluster.CompleteCurrentJob(ResizeJobStateAborted) + return errors.Wrap(err, "complete current job") +} + +func (api *API) State() string { + return api.Cluster.State() +} diff --git a/client_test.go b/client_test.go index f2bd18e50..cdbaba788 100644 --- a/client_test.go +++ b/client_test.go @@ -36,10 +36,10 @@ func createCluster(c *pilosa.Cluster) ([]*test.Server, []*test.Holder) { for i := 0; i < numNodes; i++ { hldr[i] = test.MustOpenHolder() server[i] = test.NewServer() - server[i].Handler.Cluster = c - server[i].Handler.Cluster.Nodes[i].URI = server[i].HostURI() - server[i].Handler.Holder = hldr[i].Holder - server[i].Handler.Node = server[i].Handler.Cluster.Nodes[i] + server[i].Handler.API.URI = server[i].HostURI() + server[i].Handler.API.Cluster = c + server[i].Handler.API.Cluster.Nodes[i].URI = server[i].HostURI() + server[i].Handler.API.Holder = hldr[i].Holder } return server, hldr } @@ -86,7 +86,7 @@ func TestClient_MultiNode(t *testing.T) { // Create a dispersed set of bitmaps across 3 nodes such that each individual node and slice width increment would reveal a different TopN. sliceNums := []uint64{1, 2, 6} for i, num := range sliceNums { - owns := s[i].Handler.Handler.Cluster.OwnsSlices("i", 20, s[i].HostURI()) + owns := s[i].Handler.Handler.API.Cluster.OwnsSlices("i", 20, s[i].HostURI()) ownsNum := false for _, ownNum := range owns { if ownNum == num { @@ -217,10 +217,10 @@ func TestClient_Import(t *testing.T) { s := test.NewServer() defer s.Close() - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0].URI = s.HostURI() - s.Handler.Holder = hldr.Holder - s.Handler.Node = s.Handler.Cluster.Nodes[0] + s.Handler.API.URI = s.HostURI() + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder // Send import request. c := test.MustNewClient(s.Host(), defaultClient) @@ -268,10 +268,10 @@ func TestClient_ImportInverseEnabled(t *testing.T) { s := test.NewServer() defer s.Close() - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0].URI = s.HostURI() - s.Handler.Holder = hldr.Holder - s.Handler.Node = s.Handler.Cluster.Nodes[0] + s.Handler.API.URI = s.HostURI() + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder // Send import request. c := test.MustNewClient(s.Host(), defaultClient) @@ -317,10 +317,10 @@ func TestClient_ImportValue(t *testing.T) { s := test.NewServer() defer s.Close() - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0].URI = s.HostURI() - s.Handler.Holder = hldr.Holder - s.Handler.Node = s.Handler.Cluster.Nodes[0] + s.Handler.API.URI = s.HostURI() + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder // Send import request. c := test.MustNewClient(s.Host(), defaultClient) @@ -355,10 +355,10 @@ func TestClient_BackupRestore(t *testing.T) { s := test.NewServer() defer s.Close() - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0].URI = s.HostURI() - s.Handler.Holder = hldr.Holder - s.Handler.Node = s.Handler.Cluster.Nodes[0] + s.Handler.API.URI = s.HostURI() + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder c := test.MustNewClient(s.Host(), defaultClient) @@ -420,10 +420,11 @@ func TestClient_BackupInverseView(t *testing.T) { s := test.NewServer() defer s.Close() - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0].URI = s.HostURI() - s.Handler.Holder = hldr.Holder - s.Handler.Node = s.Handler.Cluster.Nodes[0] + + s.Handler.API.URI = s.HostURI() + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder c := test.MustNewClient(s.Host(), defaultClient) @@ -457,10 +458,10 @@ func TestClient_BackupInvalidView(t *testing.T) { s := test.NewServer() defer s.Close() - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0].URI = s.HostURI() - s.Handler.Holder = hldr.Holder - s.Handler.Node = s.Handler.Cluster.Nodes[0] + s.Handler.API.URI = s.HostURI() + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder c := test.MustNewClient(s.Host(), defaultClient) @@ -486,10 +487,10 @@ func TestClient_FragmentBlocks(t *testing.T) { s := test.NewServer() defer s.Close() - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0].URI = s.HostURI() - s.Handler.Holder = hldr.Holder - s.Handler.Node = s.Handler.Cluster.Nodes[0] + s.Handler.API.URI = s.HostURI() + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder // Retrieve blocks. c := test.MustNewClient(s.Host(), defaultClient) diff --git a/cluster.go b/cluster.go index 7e53628ac..c490c58cb 100644 --- a/cluster.go +++ b/cluster.go @@ -1198,7 +1198,7 @@ func (c *Cluster) CompleteCurrentJob(state string) error { c.mu.Lock() defer c.mu.Unlock() if c.currentJob == nil { - return fmt.Errorf("no resize job currently running") + return ErrResizeNotRunning } c.currentJob.SetState(state) c.currentJob = nil @@ -1747,7 +1747,7 @@ func (c *Cluster) nodeJoin(node *Node) error { func (c *Cluster) NodeLeave(node *Node) error { // Refuse the request if this is not the coordinator. if !c.IsCoordinator() { - return fmt.Errorf("Node removal requests are only valid on the Coordinator node: %s", c.CoordinatorNode().ID) + return fmt.Errorf("node removal requests are only valid on the coordinator node: %s", c.CoordinatorNode().ID) } if c.State() != ClusterStateNormal { @@ -1761,7 +1761,7 @@ func (c *Cluster) NodeLeave(node *Node) error { // Prevent removing the coordinator node (this node). if node.ID == c.Node.ID { - return fmt.Errorf("The coordinator node cannot be removed. First, make a different node the new coordinator.") + return fmt.Errorf("coordinator cannot be removed; first, make a different node the new coordinator.") } // See if resize job can be generated diff --git a/ctl/backup_test.go b/ctl/backup_test.go index d23b1a917..d80b72475 100644 --- a/ctl/backup_test.go +++ b/ctl/backup_test.go @@ -50,12 +50,10 @@ func TestBackupCommand_Run(t *testing.T) { if err != nil { t.Fatal(err) } - node := &pilosa.Node{ID: "node", URI: *uri} - - s.Handler.Node = node - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0].URI = *uri - s.Handler.Holder = hldr.Holder + s.Handler.API.URI = *uri + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder cm := NewBackupCommand(stdin, stdout, stderr) file, err := ioutil.TempFile("", "import.csv") diff --git a/ctl/export_test.go b/ctl/export_test.go index 5d1d5a4d7..b414f711f 100644 --- a/ctl/export_test.go +++ b/ctl/export_test.go @@ -63,12 +63,10 @@ func TestExportCommand_Run(t *testing.T) { if err != nil { t.Fatal(err) } - node := &pilosa.Node{ID: "node", URI: *uri} - - s.Handler.Node = node - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0] = node - s.Handler.Holder = hldr.Holder + s.Handler.API.URI = *uri + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder cm.Host = s.Host() http.DefaultClient.Do(test.MustNewHTTPRequest("POST", s.URL+"/index/i", strings.NewReader(""))) diff --git a/ctl/import_test.go b/ctl/import_test.go index 9022efae9..bd259e342 100644 --- a/ctl/import_test.go +++ b/ctl/import_test.go @@ -69,12 +69,10 @@ func TestImportCommand_Run(t *testing.T) { if err != nil { t.Fatal(err) } - node := &pilosa.Node{ID: "node", URI: *uri} - - s.Handler.Node = node - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0] = node - s.Handler.Holder = hldr.Holder + s.Handler.API.URI = *uri + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder cm.Host = s.Host() cm.Index = "i" @@ -111,16 +109,15 @@ func TestImportCommand_RunValue(t *testing.T) { if err != nil { t.Fatal(err) } - node := &pilosa.Node{ID: "node", URI: *uri} - s.Handler.Node = node - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0] = node - s.Handler.Holder = hldr.Holder + s.Handler.API.URI = *uri + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder cm.Host = s.Host() http.DefaultClient.Do(MustNewHTTPRequest("POST", s.URL+"/index/i", strings.NewReader(""))) - http.DefaultClient.Do(MustNewHTTPRequest("POST", s.URL+"/index/i/frame/f", strings.NewReader(""))) + http.DefaultClient.Do(MustNewHTTPRequest("POST", s.URL+"/index/i/frame/f", strings.NewReader(`{"options":{"rangeEnabled": true, "fields": [{"name": "foo", "type": "int", "min": 0, "max": 100}]}}`))) cm.Index = "i" cm.Frame = "f" diff --git a/ctl/restore_test.go b/ctl/restore_test.go index bb8eb3b5f..04cb30415 100644 --- a/ctl/restore_test.go +++ b/ctl/restore_test.go @@ -52,12 +52,10 @@ func TestRestoreCommand_Run(t *testing.T) { if err != nil { t.Fatal(err) } - node := &pilosa.Node{ID: "node", URI: *uri} - - s.Handler.Node = node - s.Handler.Cluster = test.NewCluster(1) - s.Handler.Cluster.Nodes[0].URI = *uri - s.Handler.Holder = hldr.Holder + s.Handler.API.URI = *uri + s.Handler.API.Cluster = test.NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() + s.Handler.API.Holder = hldr.Holder cm := NewRestoreCommand(stdin, stdout, stderr) cm.Path = file.Name() diff --git a/executor_test.go b/executor_test.go index 1a23eba14..f2e0d210d 100644 --- a/executor_test.go +++ b/executor_test.go @@ -927,7 +927,7 @@ func TestExecutor_Execute_Remote_Bitmap(t *testing.T) { // The local node owns slice 1. hldr := test.MustOpenHolder() defer hldr.Close() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).MustSetBits(10, (1*SliceWidth)+1) e := test.NewExecutor(hldr.Holder, c) @@ -961,7 +961,7 @@ func TestExecutor_Execute_Remote_Count(t *testing.T) { // Create local executor data. The local node owns slice 1. hldr := test.MustOpenHolder() defer hldr.Close() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).MustSetBits(10, (2*SliceWidth)+1) hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).MustSetBits(10, (2*SliceWidth)+2) @@ -1004,7 +1004,7 @@ func TestExecutor_Execute_Remote_SetBit(t *testing.T) { // Create local executor data. hldr := test.MustOpenHolder() defer hldr.Close() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder // Create frame. if _, err := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{}).CreateFrame("f", pilosa.FrameOptions{}); err != nil { @@ -1056,7 +1056,7 @@ func TestExecutor_Execute_Remote_SetBit_With_Timestamp(t *testing.T) { // Create local executor data. hldr := test.MustOpenHolder() defer hldr.Close() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder // Create frame. if f, err := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{}).CreateFrame("f", pilosa.FrameOptions{}); err != nil { @@ -1130,7 +1130,7 @@ func TestExecutor_Execute_Remote_TopN(t *testing.T) { // Create local executor data on slice 2 & 4. hldr := test.MustOpenHolder() defer hldr.Close() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).MustSetBits(30, (2*SliceWidth)+1) hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 4).MustSetBits(30, (4*SliceWidth)+2) diff --git a/handler.go b/handler.go index ffd28da46..d7536f738 100644 --- a/handler.go +++ b/handler.go @@ -16,9 +16,7 @@ package pilosa import ( "context" - "encoding/csv" "encoding/json" - "errors" "expvar" "fmt" "io" @@ -27,36 +25,25 @@ import ( "net/url" // Imported for its side-effect of registering pprof endpoints with the server. _ "net/http/pprof" + "reflect" "runtime/debug" "strconv" "strings" "time" - - "reflect" + "unicode" "github.com/gogo/protobuf/proto" "github.com/gorilla/mux" "github.com/pilosa/pilosa/internal" "github.com/pilosa/pilosa/pql" - - "unicode" + "github.com/pkg/errors" ) // Handler represents an HTTP handler. type Handler struct { - Holder *Holder - Broadcaster Broadcaster - BroadcastHandler BroadcastHandler - StatusHandler StatusHandler + Router *mux.Router - FileSystem FileSystem - - // Local hostname & cluster configuration. - Node *Node - Cluster *Cluster - RemoteClient *http.Client - - Router *mux.Router + FileSystem FileSystem NormalRouter *mux.Router RestrictedRouter *mux.Router @@ -69,6 +56,8 @@ type Handler struct { // Keeps the query argument validators for each handler validators map[string]*queryValidationSpec + + API *API } // externalPrefixFlag denotes endpoints that are intended to be exposed to clients. @@ -91,9 +80,6 @@ type errorResponse struct { // NewHandler returns a new instance of Handler with a default logger. func NewHandler() *Handler { handler := &Handler{ - Broadcaster: NopBroadcaster, - //BroadcastHandler: NopBroadcastHandler, // TODO: implement the nop - //StatusHandler: NopStatusHandler, // TODO: implement the nop FileSystem: NopFileSystem, Logger: NopLogger, } @@ -229,7 +215,7 @@ func loadNormal(router *mux.Router, handler *Handler) { } func (h *Handler) reportRestricted(w http.ResponseWriter, r *http.Request) { - http.Error(w, fmt.Sprintf("not allowed in cluster state %s", h.Cluster.State()), http.StatusMethodNotAllowed) + http.Error(w, fmt.Sprintf("not allowed in cluster state %s", h.API.State()), http.StatusMethodNotAllowed) } func (h *Handler) methodNotAllowedHandler(w http.ResponseWriter, r *http.Request) { @@ -253,25 +239,25 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { dif := time.Since(t) // Calculate per request StatsD metrics when the handler is fully configured. - if h.Holder != nil && h.Cluster != nil { - statsTags := make([]string, 0, 3) + statsTags := make([]string, 0, 3) - if h.Cluster.LongQueryTime > 0 && dif > h.Cluster.LongQueryTime { - h.Logger.Printf("%s %s %v", r.Method, r.URL.String(), dif) - statsTags = append(statsTags, "slow_query") - } + longQueryTime := h.API.ClusterLongQueryTime() + if longQueryTime > 0 && dif > longQueryTime { + h.Logger.Printf("%s %s %v", r.Method, r.URL.String(), dif) + statsTags = append(statsTags, "slow_query") + } - pathParts := strings.Split(r.URL.Path, "/") - endpointName := strings.Join(pathParts, "_") + pathParts := strings.Split(r.URL.Path, "/") + endpointName := strings.Join(pathParts, "_") - if externalPrefixFlag[pathParts[1]] { - statsTags = append(statsTags, "external") - } + if externalPrefixFlag[pathParts[1]] { + statsTags = append(statsTags, "external") + } - // useragent tag identifies internal/external endpoints - statsTags = append(statsTags, "useragent:"+r.UserAgent()) - - stats := h.Holder.Stats.WithTags(statsTags...) + // useragent tag identifies internal/external endpoints + statsTags = append(statsTags, "useragent:"+r.UserAgent()) + stats := h.API.StatsWithTags(statsTags) + if stats != nil { stats.Histogram("http."+endpointName, float64(dif), 0.1) } } @@ -293,8 +279,9 @@ func (h *Handler) handleWebUI(w http.ResponseWriter, r *http.Request) { // handleGetSchema handles GET /schema requests. func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) { + schema := h.API.Schema(r.Context()) if err := json.NewEncoder(w).Encode(getSchemaResponse{ - Indexes: h.Holder.Schema(), + Indexes: schema, }); err != nil { h.Logger.Printf("write schema response error: %s", err) } @@ -302,13 +289,16 @@ 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) { - pb, err := h.StatusHandler.ClusterStatus() + pb, err := h.API.Status(r.Context()) if err != nil { h.Logger.Printf("cluster status error: %s", err) return } - cs := pb.(*internal.ClusterStatus) + cs, ok := pb.(*internal.ClusterStatus) + if !ok { + panic("status is not a status") + } if err := json.NewEncoder(w).Encode(getStatusResponse{ State: cs.State, Nodes: DecodeNodes(cs.Nodes), @@ -328,8 +318,6 @@ type getStatusResponse struct { // handlePostQuery handles /query requests. func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { - indexName := mux.Vars(r)["index"] - // Parse incoming request. req, err := h.readQueryRequest(r) if err != nil { @@ -337,48 +325,16 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { h.writeQueryResponse(w, r, &QueryResponse{Err: err}) return } + // TODO: Remove + req.Index = mux.Vars(r)["index"] - // Build execution options. - opt := &ExecOptions{ - Remote: req.Remote, - ExcludeAttrs: req.ExcludeAttrs, - ExcludeBits: req.ExcludeBits, - } - - // Parse query string. - q, err := pql.NewParser(strings.NewReader(req.Query)).Parse() + resp, err := h.API.ExecuteQuery(r.Context(), req) if err != nil { w.WriteHeader(http.StatusBadRequest) h.writeQueryResponse(w, r, &QueryResponse{Err: err}) return } - // Execute the query. - results, err := h.Executor.Execute(r.Context(), indexName, q, req.Slices, opt) - resp := &QueryResponse{Results: results, Err: err} - - // Fill column attributes if requested. - if req.ColumnAttrs && !req.ExcludeBits { - // Consolidate all column ids across all calls. - var columnIDs []uint64 - for _, result := range results { - bm, ok := result.(*Bitmap) - if !ok { - continue - } - columnIDs = uint64Slice(columnIDs).merge(bm.Bits()) - } - - // Retrieve column attributes across all calls. - columnAttrSets, err := h.readColumnAttrSets(h.Holder.Index(indexName), columnIDs) - if err != nil { - w.WriteHeader(http.StatusInternalServerError) - h.writeQueryResponse(w, r, &QueryResponse{Err: err}) - return - } - resp.ColumnAttrSets = columnAttrSets - } - // Set appropriate status code, if there is an error. if resp.Err != nil { switch resp.Err { @@ -390,7 +346,7 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { } // Write response back to client. - if err := h.writeQueryResponse(w, r, resp); err != nil { + if err := h.writeQueryResponse(w, r, &resp); err != nil { h.Logger.Printf("write query response error: %s", err) } } @@ -398,8 +354,8 @@ 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 err := json.NewEncoder(w).Encode(getSlicesMaxResponse{ - Standard: h.Holder.MaxSlices(), - Inverse: h.Holder.MaxInverseSlices(), + Standard: h.API.MaxSlices(r.Context()), + Inverse: h.API.MaxInverseSlices(r.Context()), }); err != nil { h.Logger.Printf("write slices-max response error: %s", err) } @@ -418,9 +374,9 @@ func (h *Handler) handleGetIndexes(w http.ResponseWriter, r *http.Request) { // handleGetIndex handles GET /index/ requests. func (h *Handler) handleGetIndex(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) + index, err := h.API.ReadIndex(r.Context(), indexName) + if err != nil { + http.Error(w, err.Error(), http.StatusNotFound) return } @@ -502,28 +458,17 @@ type postIndexResponse struct{} // handleDeleteIndex handles DELETE /index request. func (h *Handler) handleDeleteIndex(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] - - // Delete index from the holder. - if err := h.Holder.DeleteIndex(indexName); err != nil { + err := h.API.DeleteIndex(r.Context(), indexName) + if err != nil { + h.Logger.Printf("problem deleting index: %s", err) http.Error(w, err.Error(), http.StatusInternalServerError) return } - // Send the delete index message to all nodes. - err := h.Broadcaster.SendSync( - &internal.DeleteIndexMessage{ - Index: indexName, - }) - if err != nil { - h.Logger.Printf("problem sending DeleteIndex message: %s", err) - } - // Encode response. if err := json.NewEncoder(w).Encode(deleteIndexResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) } - - h.Holder.Stats.Count("deleteIndex", 1, 1.0) } type deleteIndexResponse struct{} @@ -543,8 +488,7 @@ func (h *Handler) handlePostIndex(w http.ResponseWriter, r *http.Request) { return } - // Create index. - _, err = h.Holder.CreateIndex(indexName, req.Options) + _, err = h.API.CreateIndex(r.Context(), indexName, req.Options) if err == ErrIndexExists { http.Error(w, err.Error(), http.StatusConflict) return @@ -553,24 +497,10 @@ func (h *Handler) handlePostIndex(w http.ResponseWriter, r *http.Request) { return } - // Send the create index message to all nodes. - err = h.Broadcaster.SendSync( - &internal.CreateIndexMessage{ - Index: indexName, - Meta: req.Options.Encode(), - }) - if err != nil { - h.Logger.Printf("problem sending CreateIndex message: %s", err) - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - // Encode response. if err := json.NewEncoder(w).Encode(postIndexResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) } - - h.Holder.Stats.Count("createIndex", 1, 1.0) } // handlePatchIndexTimeQuantum handles PATCH /index/time_quantum request. @@ -591,16 +521,12 @@ func (h *Handler) handlePatchIndexTimeQuantum(w http.ResponseWriter, r *http.Req return } - // Retrieve index by name. - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - // Set default time quantum on index. - if err := index.SetTimeQuantum(tq); err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) + if err = h.API.ModifyIndexTimeQuantum(r.Context(), indexName, tq); err != nil { + if err == ErrIndexNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } @@ -627,34 +553,14 @@ func (h *Handler) handlePostIndexAttrDiff(w http.ResponseWriter, r *http.Request return } - // Retrieve index from holder. - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - // Retrieve local blocks. - blks, err := index.ColumnAttrStore().Blocks() + attrs, err := h.API.IndexAttrDiff(r.Context(), indexName, req.Blocks) if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Read all attributes from all mismatched blocks. - attrs := make(map[uint64]map[string]interface{}) - for _, blockID := range AttrBlocks(blks).Diff(req.Blocks) { - // Retrieve block data. - m, err := index.ColumnAttrStore().BlockData(blockID) - if err != nil { + if err == ErrIndexNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Copy to index-wide struct. - for k, v := range m { - attrs[k] = v } + return } // Encode response. @@ -688,41 +594,22 @@ func (h *Handler) handlePostFrame(w http.ResponseWriter, r *http.Request) { http.Error(w, err.Error(), http.StatusBadRequest) return } - - // Find index. - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - // Create frame. - _, err = index.CreateFrame(frameName, req.Options) - if err == ErrFrameExists { - http.Error(w, err.Error(), http.StatusConflict) - return - } else if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Send the create frame message to all nodes. - err = h.Broadcaster.SendSync( - &internal.CreateFrameMessage{ - Index: indexName, - Frame: frameName, - Meta: req.Options.Encode(), - }) + _, err = h.API.CreateFrame(r.Context(), indexName, frameName, req.Options) if err != nil { - h.Logger.Printf("problem sending CreateFrame message: %s", err) + switch err { + case ErrIndexNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + case ErrFrameExists: + http.Error(w, err.Error(), http.StatusConflict) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) + } + return } - // Encode response. if err := json.NewEncoder(w).Encode(postFrameResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) } - - h.Holder.Stats.CountWithCustomTags("createFrame", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) } type _postFrameRequest postFrameRequest @@ -775,37 +662,22 @@ func (h *Handler) handleDeleteFrame(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] frameName := mux.Vars(r)["frame"] - // Find index. - index := h.Holder.Index(indexName) - if index == nil { - if err := json.NewEncoder(w).Encode(deleteIndexResponse{}); err != nil { - h.Logger.Printf("response encoding error: %s", err) + err := h.API.DeleteFrame(r.Context(), indexName, frameName) + if err != nil { + if err == ErrIndexNotFound { + if err := json.NewEncoder(w).Encode(deleteIndexResponse{}); err != nil { + h.Logger.Printf("response encoding error: %s", err) + } + return } - return - } - - // Delete frame from the index. - if err := index.DeleteFrame(frameName); err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) return } - // Send the delete frame message to all nodes. - err := h.Broadcaster.SendSync( - &internal.DeleteFrameMessage{ - Index: indexName, - Frame: frameName, - }) - if err != nil { - h.Logger.Printf("problem sending DeleteFrame message: %s", err) - } - // Encode response. if err := json.NewEncoder(w).Encode(deleteFrameResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) } - - h.Holder.Stats.CountWithCustomTags("deleteFrame", 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) } type deleteFrameResponse struct{} @@ -829,16 +701,12 @@ func (h *Handler) handlePatchFrameTimeQuantum(w http.ResponseWriter, r *http.Req return } - // Retrieve index by name. - f := h.Holder.Frame(indexName, frameName) - if f == nil { - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - // Set default time quantum on index. - if err := f.SetTimeQuantum(tq); err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) + if err := h.API.ModifyFrameTimeQuantum(r.Context(), indexName, frameName, tq); err != nil { + if err == ErrFragmentNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } @@ -867,13 +735,6 @@ func (h *Handler) handlePostFrameField(w http.ResponseWriter, r *http.Request) { return } - // Retrieve frame by name. - f := h.Holder.Frame(indexName, frameName) - if f == nil { - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - field := &Field{ Name: fieldName, Type: req.Type, @@ -881,23 +742,15 @@ func (h *Handler) handlePostFrameField(w http.ResponseWriter, r *http.Request) { Max: req.Max, } - // Create new field. - if err := f.CreateField(field); err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) + if err := h.API.CreateFrameField(r.Context(), indexName, frameName, field); err != nil { + if err == ErrFrameNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } - // Send the create field message to all nodes. - err := h.Broadcaster.SendSync( - &internal.CreateFieldMessage{ - Index: indexName, - Frame: frameName, - Field: encodeField(field), - }) - if err != nil { - h.Logger.Printf("problem sending CreateField message: %s", err) - } - // Encode response. if err := json.NewEncoder(w).Encode(postFrameFieldResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) @@ -918,30 +771,15 @@ func (h *Handler) handleDeleteFrameField(w http.ResponseWriter, r *http.Request) frameName := mux.Vars(r)["frame"] fieldName := mux.Vars(r)["field"] - // Retrieve frame by name. - f := h.Holder.Frame(indexName, frameName) - if f == nil { - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) + if err := h.API.DeleteFrameField(r.Context(), indexName, frameName, fieldName); err != nil { + if err == ErrFrameNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } - // Delete field. - if err := f.DeleteField(fieldName); err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Send the delete field message to all nodes. - err := h.Broadcaster.SendSync( - &internal.DeleteFieldMessage{ - Index: indexName, - Frame: frameName, - Field: fieldName, - }) - if err != nil { - h.Logger.Printf("problem sending DeleteField message: %s", err) - } - // Encode response. if err := json.NewEncoder(w).Encode(deleteFrameFieldResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) @@ -952,24 +790,18 @@ func (h *Handler) handleGetFrameFields(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] frameName := mux.Vars(r)["frame"] - index := h.Holder.index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - frame := index.frame(frameName) - if frame == nil { - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - fields, err := frame.GetFields() - if err == ErrFrameFieldsNotAllowed { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } else if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) + fields, err := h.API.FrameFields(r.Context(), indexName, frameName) + if err != nil { + switch err { + case ErrIndexNotFound: + fallthrough + case ErrFrameNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + case ErrFrameFieldsNotAllowed: + http.Error(w, err.Error(), http.StatusBadRequest) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } @@ -992,15 +824,16 @@ func (h *Handler) handleGetFrameViews(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] frameName := mux.Vars(r)["frame"] - // Retrieve views. - f := h.Holder.Frame(indexName, frameName) - if f == nil { - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) + views, err := h.API.FrameViews(r.Context(), indexName, frameName) + if err != nil { + if err == ErrFrameNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } - // Fetch views. - views := f.Views() names := make([]string, len(views)) for i := range views { names[i] = views[i].Name() @@ -1018,31 +851,13 @@ func (h *Handler) handleDeleteView(w http.ResponseWriter, r *http.Request) { frameName := mux.Vars(r)["frame"] viewName := mux.Vars(r)["view"] - // Retrieve frame. - f := h.Holder.Frame(indexName, frameName) - if f == nil { - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - // Delete the view. - if err := f.DeleteView(viewName); err != nil { - // Ingore this error because views do not exist on all nodes due to slice distribution. - if err != ErrInvalidView { + if err := h.API.DeleteView(r.Context(), indexName, frameName, viewName); err != nil { + if err == ErrFrameNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { http.Error(w, err.Error(), http.StatusBadRequest) - return } - } - - // Send the delete view message to all nodes. - err := h.Broadcaster.SendSync( - &internal.DeleteViewMessage{ - Index: indexName, - Frame: frameName, - View: viewName, - }) - if err != nil { - h.Logger.Printf("problem sending DeleteView message: %s", err) + return } // Encode response. @@ -1069,34 +884,15 @@ func (h *Handler) handlePostFrameAttrDiff(w http.ResponseWriter, r *http.Request return } - // Retrieve index from holder. - f := h.Holder.Frame(indexName, frameName) - if f == nil { - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - // Retrieve local blocks. - blks, err := f.RowAttrStore().Blocks() + attrs, err := h.API.FrameAttrDiff(r.Context(), indexName, frameName, req.Blocks) if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Read all attributes from all mismatched blocks. - attrs := make(map[uint64]map[string]interface{}) - for _, blockID := range AttrBlocks(blks).Diff(req.Blocks) { - // Retrieve block data. - m, err := f.RowAttrStore().BlockData(blockID) - if err != nil { + switch err { + case ErrFragmentNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + default: http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Copy to index-wide struct. - for k, v := range m { - attrs[k] = v } + return } // Encode response. @@ -1254,44 +1050,17 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) { return } - // Convert timestamps to time.Time. - timestamps := make([]*time.Time, len(req.Timestamps)) - for i, ts := range req.Timestamps { - if ts == 0 { - continue + if err := h.API.Import(r.Context(), req); err != nil { + switch err { + case ErrIndexNotFound: + fallthrough + case ErrFrameNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + case ErrClusterDoesNotOwnSlice: + http.Error(w, err.Error(), http.StatusPreconditionFailed) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) } - t := time.Unix(0, ts) - timestamps[i] = &t - } - - // Validate that this handler owns the slice. - if !h.Cluster.OwnsFragment(h.Node.ID, req.Index, req.Slice) { - msg := fmt.Sprintf("host does not own slice %s-%s slice:%d", h.Node.ID, req.Index, req.Slice) - http.Error(w, msg, http.StatusPreconditionFailed) - return - } - - // Find the Index. - h.Logger.Printf("importing: %s %s %d", req.Index, req.Frame, req.Slice) - index := h.Holder.Index(req.Index) - if index == nil { - h.Logger.Printf("fragment error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrIndexNotFound.Error()) - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - // Retrieve frame. - f := index.Frame(req.Frame) - if f == nil { - h.Logger.Printf("frame error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrFrameNotFound.Error()) - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - // Import into fragment. - err = f.Import(req.RowIDs, req.ColumnIDs, timestamps) - if err != nil { - h.Logger.Printf("import error: index=%s, frame=%s, slice=%d, bits=%d, err=%s", req.Index, req.Frame, req.Slice, len(req.ColumnIDs), err) return } @@ -1334,34 +1103,17 @@ func (h *Handler) handlePostImportValue(w http.ResponseWriter, r *http.Request) return } - // Validate that this handler owns the slice. - if !h.Cluster.OwnsFragment(h.Node.ID, req.Index, req.Slice) { - msg := fmt.Sprintf("host does not own slice %s-%s slice:%d", h.Node.ID, req.Index, req.Slice) - http.Error(w, msg, http.StatusPreconditionFailed) - return - } - - // Find the Index. - h.Logger.Printf("importing: %s %s %d", req.Index, req.Frame, req.Slice) - index := h.Holder.Index(req.Index) - if index == nil { - h.Logger.Printf("fragment error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrIndexNotFound.Error()) - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - // Retrieve frame. - f := index.Frame(req.Frame) - if f == nil { - h.Logger.Printf("frame error: index=%s, frame=%s, slice=%d, err=%s", req.Index, req.Frame, req.Slice, ErrFrameNotFound.Error()) - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - // Import into fragment. - err = f.ImportValue(req.Field, req.ColumnIDs, req.Values) - if err != nil { - h.Logger.Printf("import error: index=%s, frame=%s, slice=%d, field=%s, bits=%d, err=%s", req.Index, req.Frame, req.Slice, req.Field, len(req.ColumnIDs), err) + if err = h.API.ImportValue(r.Context(), req); err != nil { + switch err { + case ErrIndexNotFound: + fallthrough + case ErrFrameNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + case ErrClusterDoesNotOwnSlice: + http.Error(w, err.Error(), http.StatusPreconditionFailed) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } @@ -1400,35 +1152,17 @@ func (h *Handler) handleGetExportCSV(w http.ResponseWriter, r *http.Request) { return } - // Validate that this handler owns the slice. - if !h.Cluster.OwnsFragment(h.Node.ID, index, slice) { - msg := fmt.Sprintf("host does not own slice %s-%s slice:%d", h.Node.ID, index, slice) - http.Error(w, msg, http.StatusPreconditionFailed) + if err = h.API.ExportCSV(r.Context(), index, frame, view, slice, w); err != nil { + switch err { + case ErrFragmentNotFound: + break + case ErrClusterDoesNotOwnSlice: + http.Error(w, err.Error(), http.StatusPreconditionFailed) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } - - // Find the fragment. - f := h.Holder.Fragment(index, frame, view, slice) - if f == nil { - return - } - - // Wrap writer with a CSV writer. - cw := csv.NewWriter(w) - - // Iterate over each bit. - if err := f.ForEachBit(func(rowID, columnID uint64) error { - return cw.Write([]string{ - strconv.FormatUint(rowID, 10), - strconv.FormatUint(columnID, 10), - }) - }); err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Ensure data is flushed. - cw.Flush() } // handleGetFragmentNodes handles /fragment/nodes requests. @@ -1444,7 +1178,7 @@ func (h *Handler) handleGetFragmentNodes(w http.ResponseWriter, r *http.Request) } // Retrieve fragment owner nodes. - nodes := h.Cluster.FragmentNodes(index, slice) + nodes := h.API.FragmentNodes(r.Context(), index, slice) // Write to response. if err := json.NewEncoder(w).Encode(nodes); err != nil { @@ -1463,9 +1197,9 @@ func (h *Handler) handleGetFragmentData(w http.ResponseWriter, r *http.Request) } // Retrieve fragment from holder. - f := h.Holder.Fragment(q.Get("index"), q.Get("frame"), q.Get("view"), slice) - if f == nil { - http.Error(w, "fragment not found", http.StatusNotFound) + f, err := h.API.FragmentData(r.Context(), q.Get("index"), q.Get("frame"), q.Get("view"), slice) + if err != nil { + http.Error(w, err.Error(), http.StatusNotFound) return } @@ -1485,31 +1219,12 @@ func (h *Handler) handlePostFragmentData(w http.ResponseWriter, r *http.Request) return } - // Retrieve frame. - f := h.Holder.Frame(q.Get("index"), q.Get("frame")) - if f == nil { - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - // Retrieve view. - view, err := f.CreateViewIfNotExists(q.Get("view")) - if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Retrieve fragment from frame. - frag, err := view.CreateFragmentIfNotExists(slice) - if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Read fragment in from request body. - if _, err := frag.ReadFrom(r.Body); err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return + if err = h.API.WriteFragmentData(r.Context(), q.Get("index"), q.Get("frame"), q.Get("view"), slice, r.Body); err != nil { + if err == ErrFrameNotFound { + http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) + } else { + http.Error(w, err.Error(), http.StatusInternalServerError) + } } } @@ -1525,19 +1240,16 @@ func (h *Handler) handleGetFragmentBlockData(w http.ResponseWriter, r *http.Requ return } - // Retrieve fragment from holder. - f := h.Holder.Fragment(req.Index, req.Frame, req.View, req.Slice) - if f == nil { - http.Error(w, ErrFragmentNotFound.Error(), http.StatusNotFound) + resp, err := h.API.FragmentBlockData(r.Context(), req) + if err != nil { + if err == ErrFragmentNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } - // Read data - var resp internal.BlockDataResponse - if f != nil { - resp.RowIDs, resp.ColumnIDs = f.BlockData(int(req.Block)) - } - // Encode response. buf, err := proto.Marshal(&resp) if err != nil { @@ -1561,16 +1273,16 @@ func (h *Handler) handleGetFragmentBlocks(w http.ResponseWriter, r *http.Request return } - // Retrieve fragment from holder. - f := h.Holder.Fragment(q.Get("index"), q.Get("frame"), q.Get("view"), slice) - if f == nil { - http.Error(w, "fragment not found", http.StatusNotFound) + blocks, err := h.API.FragmentBlocks(r.Context(), q.Get("index"), q.Get("frame"), q.Get("view"), slice) + if err != nil { + if err == ErrFragmentNotFound { + http.Error(w, err.Error(), http.StatusNotFound) + } else { + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } - // Retrieve blocks. - blocks := f.Blocks() - // Encode response. if err := json.NewEncoder(w).Encode(getFragmentBlocksResponse{ Blocks: blocks, @@ -1602,80 +1314,23 @@ func (h *Handler) handlePostFrameRestore(w http.ResponseWriter, r *http.Request) http.Error(w, err.Error(), http.StatusBadRequest) } - // Create a client for the remote cluster. - client := NewInternalHTTPClientFromURI(host, h.RemoteClient) - - // Determine the maximum number of slices. - maxSlices, err := client.MaxSliceByIndex(r.Context()) - if err != nil { - http.Error(w, "cannot determine remote slice count: "+err.Error(), http.StatusInternalServerError) - return - } - - // Retrieve frame. - f := h.Holder.Frame(indexName, frameName) - if f == nil { - http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound) - return - } - - // Retrieve list of all views. - views, err := client.FrameViews(r.Context(), indexName, frameName) - if err != nil { - http.Error(w, "cannot retrieve frame views: "+err.Error(), http.StatusInternalServerError) - return - } - - // Loop over each slice and import it if this node owns it. - for slice := uint64(0); slice <= maxSlices[indexName]; slice++ { - // Ignore this slice if we don't own it. - if !h.Cluster.OwnsFragment(h.Node.ID, indexName, slice) { - continue - } - - // Loop over view names. - for _, view := range views { - // Create view. - v, err := f.CreateViewIfNotExists(view) - if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Otherwise retrieve the local fragment. - frag, err := v.CreateFragmentIfNotExists(slice) - if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - // Stream backup from remote node. - rd, err := client.BackupSlice(r.Context(), indexName, frameName, view, slice) - if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } else if rd == nil { - continue // slice doesn't exist - } - - // Restore to local frame and always close reader. - if err := func() error { - defer rd.Close() - if _, err := frag.ReadFrom(rd); err != nil { - return err - } - return nil - }(); err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - } + err = h.API.RestoreFrame(r.Context(), indexName, frameName, host) + switch err { + case nil: + break + case ErrFrameNotFound: + fallthrough + case ErrFragmentNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) } } // handleGetHosts handles /hosts requests. func (h *Handler) handleGetHosts(w http.ResponseWriter, r *http.Request) { - if err := json.NewEncoder(w).Encode(h.Cluster.Nodes); err != nil { + hosts := h.API.ClusterHosts(r.Context()) + if err := json.NewEncoder(w).Encode(hosts); err != nil { h.Logger.Printf("write version response error: %s", err) } } @@ -1859,13 +1514,6 @@ func (h *Handler) handlePostInputDefinition(w http.ResponseWriter, r *http.Reque indexName := mux.Vars(r)["index"] inputDefName := mux.Vars(r)["input-definition"] - // Find index. - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - // Decode request. var req InputDefinitionInfo err := json.NewDecoder(r.Body).Decode(&req) @@ -1874,34 +1522,28 @@ func (h *Handler) handlePostInputDefinition(w http.ResponseWriter, r *http.Reque return } - if err := req.Validate(); err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) + if err = h.API.CreateInputDefinition(r.Context(), indexName, inputDefName, req); err != nil { + switch err { + case ErrIndexNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + case ErrInputDefinitionExists: + http.Error(w, err.Error(), http.StatusConflict) + case ErrInputDefinitionAttrsRequired: + fallthrough + case ErrInputDefinitionNameRequired: + fallthrough + case ErrInputDefinitionActionRequired: + fallthrough + case ErrInputDefinitionHasPrimaryKey: + fallthrough + case ErrInputDefinitionDupePrimaryKey: + http.Error(w, err.Error(), http.StatusBadRequest) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } - // Encode InputDefinition to its internal representation. - def := req.Encode() - def.Name = inputDefName - - // Create InputDefinition. - _, err = index.CreateInputDefinition(def) - if err == ErrInputDefinitionExists { - http.Error(w, err.Error(), http.StatusConflict) - return - } else if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } - - err = h.Broadcaster.SendSync( - &internal.CreateInputDefinitionMessage{ - Index: indexName, - Definition: def, - }) - if err != nil { - h.Logger.Printf("problem sending CreateInputDefinition message: %s", err) - } - if err := json.NewEncoder(w).Encode(defaultInputDefinitionResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) } @@ -1912,16 +1554,18 @@ func (h *Handler) handleGetInputDefinition(w http.ResponseWriter, r *http.Reques indexName := mux.Vars(r)["index"] inputDefName := mux.Vars(r)["input-definition"] - // Find index. - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - - inputDef, err := index.InputDefinition(inputDefName) + inputDef, err := h.API.InputDefinition(r.Context(), indexName, inputDefName) if err != nil { - http.Error(w, err.Error(), http.StatusNotFound) + switch err { + case nil: + break + case ErrIndexNotFound: + fallthrough + case ErrInputDefinitionNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } @@ -1931,7 +1575,6 @@ func (h *Handler) handleGetInputDefinition(w http.ResponseWriter, r *http.Reques }); err != nil { h.Logger.Printf("write status response error: %s", err) } - } // handleDeleteInputDefinition handles DELETE /input-definition request. @@ -1939,28 +1582,20 @@ func (h *Handler) handleDeleteInputDefinition(w http.ResponseWriter, r *http.Req indexName := mux.Vars(r)["index"] inputDefName := mux.Vars(r)["input-definition"] - // Find index. - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) + if err := h.API.DeleteInputDefinition(r.Context(), indexName, inputDefName); err != nil { + switch err { + case nil: + break + case ErrIndexNotFound: + fallthrough + case ErrInputDefinitionNotFound: + http.Error(w, err.Error(), http.StatusNotFound) + default: + http.Error(w, err.Error(), http.StatusNotFound) + } return } - // Delete input definition from the index. - if err := index.DeleteInputDefinition(inputDefName); err != nil { - http.Error(w, err.Error(), http.StatusNotFound) - return - } - - err := h.Broadcaster.SendSync( - &internal.DeleteInputDefinitionMessage{ - Index: indexName, - Name: inputDefName, - }) - if err != nil { - h.Logger.Printf("problem sending DeleteInputDefinition message: %s", err) - } - if err := json.NewEncoder(w).Encode(defaultInputDefinitionResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) } @@ -1972,13 +1607,6 @@ func (h *Handler) handlePostInput(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] inputDefName := mux.Vars(r)["input-definition"] - // Find index. - index := h.Holder.Index(indexName) - if index == nil { - http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound) - return - } - // Decode request. var reqs []interface{} err := json.NewDecoder(r.Body).Decode(&reqs) @@ -1986,67 +1614,44 @@ func (h *Handler) handlePostInput(w http.ResponseWriter, r *http.Request) { http.Error(w, err.Error(), http.StatusBadRequest) return } - for _, req := range reqs { - bits, err := h.InputJSONDataParser(req.(map[string]interface{}), index, inputDefName) - if err == ErrInputDefinitionNotFound { + + if err = h.API.WriteInput(r.Context(), indexName, inputDefName, reqs); err != nil { + switch err { + case nil: + break + case ErrIndexNotFound: + fallthrough + case ErrInputDefinitionNotFound: http.Error(w, err.Error(), http.StatusNotFound) - return - } else if err != nil { + default: http.Error(w, err.Error(), http.StatusBadRequest) - return - } - for fr, bs := range bits { - err := index.InputBits(fr, bs) - if err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } } + return } + if err := json.NewEncoder(w).Encode(defaultInputDefinitionResponse{}); err != nil { h.Logger.Printf("response encoding error: %s", err) } } -// handlePostClusterResizeSetCoordinator handles POST /cluster/resize/set-coordinator request. func (h *Handler) handlePostClusterResizeSetCoordinator(w http.ResponseWriter, r *http.Request) { // Decode request. var req setCoordinatorRequest err := json.NewDecoder(r.Body).Decode(&req) if err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) + http.Error(w, "decoding request "+err.Error(), http.StatusBadRequest) return } - oldNode := h.Cluster.nodeByID(h.Cluster.Coordinator) - newNode := h.Cluster.nodeByID(req.ID) - if newNode == nil { - http.Error(w, "Node with provided ID does not exist", http.StatusBadRequest) - return - } - - if err := func() error { - // If the new coordinator is this node, do the SetCoordinator directly. - if newNode.ID == h.Node.ID { - return h.Cluster.SetCoordinator(newNode) + oldNode, newNode, err := h.API.SetCoordinator(r.Context(), req.ID) + if err != nil { + if errors.Cause(err) == ErrNodeIDNotExists { + http.Error(w, "setting new coordinator: "+err.Error(), http.StatusNotFound) + } else { + http.Error(w, "setting new coordinator: "+err.Error(), http.StatusInternalServerError) } - - // Send the set-coordinator message to new node. - err := h.Broadcaster.SendTo( - newNode, - &internal.SetCoordinatorMessage{ - New: EncodeNode(newNode), - }) - if err != nil { - return fmt.Errorf("problem sending SetCoordinator message: %s", err) - } - - return nil - }(); err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) return } - // Encode response. if err := json.NewEncoder(w).Encode(setCoordinatorResponse{ Old: oldNode, @@ -2075,16 +1680,13 @@ func (h *Handler) handlePostClusterResizeRemoveNode(w http.ResponseWriter, r *ht return } - removeNode := h.Cluster.nodeByID(req.ID) - if removeNode == nil { - http.Error(w, fmt.Sprintf("Node is not a member of the cluster: %s", req.ID), http.StatusBadRequest) - return - } - - // Start the resize process (similar to NodeJoin) - err = h.Cluster.NodeLeave(removeNode) + removeNode, err := h.API.RemoveNode(req.ID) if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) + if errors.Cause(err) == ErrNodeIDNotExists { + http.Error(w, "removing node: "+err.Error(), http.StatusNotFound) + } else { + http.Error(w, "removing node: "+err.Error(), http.StatusInternalServerError) + } return } @@ -2106,21 +1708,20 @@ type removeNodeResponse struct { // handlePostClusterResizeAbort handles POST /cluster/resize/abort request. func (h *Handler) handlePostClusterResizeAbort(w http.ResponseWriter, r *http.Request) { + err := h.API.ResizeAbort() var msg string - - if err := func() error { - if !h.Cluster.IsCoordinator() { - return fmt.Errorf("abort requests must be made on the coordinator node") + if err != nil { + switch errors.Cause(err) { + case ErrNodeNotCoordinator: + http.Error(w, err.Error(), http.StatusBadRequest) + return + case ErrResizeNotRunning: + msg = err.Error() + default: + http.Error(w, err.Error(), http.StatusInternalServerError) + return } - err := h.Cluster.CompleteCurrentJob(ResizeJobStateAborted) - if err != nil { - return err - } - return nil - }(); err != nil { - msg = err.Error() } - // Encode response. if err := json.NewEncoder(w).Encode(clusterResizeAbortResponse{ Info: msg, @@ -2133,83 +1734,13 @@ type clusterResizeAbortResponse struct { Info string `json:"info"` } -// InputJSONDataParser validates input json file and executes SetBit. -func (h *Handler) InputJSONDataParser(req map[string]interface{}, index *Index, name string) (map[string][]*Bit, error) { - inputDef, err := index.InputDefinition(name) - if err != nil { - return nil, err - } - // If field in input data is not in defined definition, return error. - var colValue uint64 - validFields := make(map[string]bool) - timestampFrame := make(map[string]int64) - for _, field := range inputDef.Fields() { - validFields[field.Name] = true - if field.PrimaryKey { - value, ok := req[field.Name] - if !ok { - return nil, fmt.Errorf("primary key does not exist") - } - rawValue, ok := value.(float64) // The default JSON marshalling will interpret this as a float - if !ok { - return nil, fmt.Errorf("float64 require, got value:%s, type: %s", value, reflect.TypeOf(value)) - } - colValue = uint64(rawValue) - } - // Find frame that need to add timestamp. - for _, action := range field.Actions { - if action.ValueDestination == InputSetTimestamp { - timestampFrame[action.Frame], err = GetTimeStamp(req, field.Name) - if err != nil { - return nil, err - } - } - } - } - - for key := range req { - _, ok := validFields[key] - if !ok { - return nil, fmt.Errorf("field not found: %s", key) - } - } - - setBits := make(map[string][]*Bit) - - for _, field := range inputDef.Fields() { - // skip field that defined in definition but not in input data - if _, ok := req[field.Name]; !ok { - continue - } - - // Looking into timestampFrame map and set timestamp to the whole frame - for _, action := range field.Actions { - frame := action.Frame - timestamp := timestampFrame[action.Frame] - // Skip input data field values that are set to null - if req[field.Name] == nil { - continue - } - bit, err := HandleAction(action, req[field.Name], colValue, timestamp) - if err != nil { - return nil, fmt.Errorf("error handling action: %s, err: %s", action.ValueDestination, err) - } - if bit != nil { - setBits[frame] = append(setBits[frame], bit) - } - } - } - return setBits, nil -} - func (h *Handler) handleRecalculateCaches(w http.ResponseWriter, r *http.Request) { - err := h.Broadcaster.SendSync(&internal.RecalculateCaches{}) + err := h.API.RecalculateCaches(r.Context()) if err != nil { - w.WriteHeader(http.StatusInternalServerError) - h.writeQueryResponse(w, r, &QueryResponse{Err: err}) + http.Error(w, "recalculating caches: "+err.Error(), http.StatusInternalServerError) return } - h.Holder.RecalculateCaches() + w.WriteHeader(http.StatusNoContent) } @@ -2254,9 +1785,7 @@ func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Reques return } - // Forward the error message. - err = h.BroadcastHandler.ReceiveMessage(pb) - if err != nil { + if err := h.API.PostClusterMessage(r.Context(), pb); err != nil { http.Error(w, err.Error(), http.StatusBadRequest) return } @@ -2267,7 +1796,7 @@ func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Reques } func (h *Handler) handleGetID(w http.ResponseWriter, r *http.Request) { - _, err := w.Write([]byte(h.Cluster.Node.ID)) + _, err := w.Write([]byte(h.API.LocalID())) if err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) } @@ -2302,12 +1831,12 @@ func (s *queryValidationSpec) Optional(args ...string) *queryValidationSpec { func (s queryValidationSpec) validate(query url.Values) error { for _, req := range s.required { if query.Get(req) == "" { - return errors.New(fmt.Sprintf("%s is required", req)) + return errors.Errorf("%s is required", req) } } - for k, _ := range query { + for k := range query { if _, ok := s.args[k]; !ok { - return errors.New(fmt.Sprintf("%s is not a valid argument", k)) + return errors.Errorf("%s is not a valid argument", k) } } return nil diff --git a/handler_test.go b/handler_test.go index ba9480544..b605a3650 100644 --- a/handler_test.go +++ b/handler_test.go @@ -65,8 +65,8 @@ func TestHandler_NotFound(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/no_such_path", nil)) @@ -100,8 +100,8 @@ func TestHandler_Schema(t *testing.T) { } h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/schema", nil)) if w.Code != http.StatusOK { @@ -139,9 +139,9 @@ func TestHandler_Status(t *testing.T) { } h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) - h.StatusHandler = s + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.StatusHandler = s s.Handler = h w := httptest.NewRecorder() @@ -158,14 +158,15 @@ func TestHandler_ClusterResizeAbort(t *testing.T) { t.Run("No resize job", func(t *testing.T) { h := test.NewHandler() - h.Cluster = test.NewCluster(1) + h.API.Cluster = test.NewCluster(1) h.SetRestricted() w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/cluster/resize/abort", nil)) if w.Code != http.StatusOK { - t.Fatalf("unexpected status code: %d", w.Code) - } else if body := w.Body.String(); body != `{"info":"no resize job currently running"}`+"\n" { + bod, err := ioutil.ReadAll(w.Body) + t.Fatalf("unexpected status code: %d, bod: %s, readerr: %v", w.Code, bod, err) + } else if body := w.Body.String(); body != `{"info":"complete current job: no resize job currently running"}`+"\n" { t.Fatalf("unexpected body: %s", body) } }) @@ -186,8 +187,8 @@ func TestHandler_MaxSlices(t *testing.T) { hldr.MustCreateFragmentIfNotExists("i1", "f1", pilosa.ViewStandard, 0).MustSetBits(40, (0*SliceWidth)+8) h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/slices/max", nil)) if w.Code != http.StatusOK { @@ -227,8 +228,8 @@ func TestHandler_MaxSlices_Inverse(t *testing.T) { } h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/slices/max?inverse=true", nil)) if w.Code != http.StatusOK { @@ -244,8 +245,8 @@ func TestHandler_Query_Args_URL(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { if index != "idx0" { t.Fatalf("unexpected index: %s", index) @@ -272,8 +273,8 @@ func TestHandler_Query_Args_Protobuf(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { if index != "idx0" { t.Fatalf("unexpected index: %s", index) @@ -312,8 +313,8 @@ func TestHandler_Query_Args_Err(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/idx0/query?slices=a,b", strings.NewReader("Bitmap(id=100)"))) if w.Code != http.StatusBadRequest { @@ -339,8 +340,8 @@ func TestHandler_Query_Uint64_JSON(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return []interface{}{uint64(100)}, nil } @@ -360,8 +361,8 @@ func TestHandler_Query_Uint64_Protobuf(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return []interface{}{uint64(100)}, nil } @@ -390,8 +391,8 @@ func TestHandler_Query_Bitmap_JSON(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { bm := pilosa.NewBitmap(1, 3, 66, pilosa.SliceWidth+1) bm.Attrs = map[string]interface{}{"a": "b", "c": 1, "d": true} @@ -423,8 +424,8 @@ func TestHandler_Query_Bitmap_ColumnAttrs_JSON(t *testing.T) { } h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { bm := pilosa.NewBitmap(1, 3, 66, pilosa.SliceWidth+1) bm.Attrs = map[string]interface{}{"a": "b", "c": 1, "d": true} @@ -446,8 +447,8 @@ func TestHandler_Query_Bitmap_Protobuf(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { bm := pilosa.NewBitmap(1, pilosa.SliceWidth+1) bm.Attrs = map[string]interface{}{"a": "b", "c": int64(1), "d": true} @@ -494,8 +495,8 @@ func TestHandler_Query_Bitmap_ColumnAttrs_Protobuf(t *testing.T) { } h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { bm := pilosa.NewBitmap(1, pilosa.SliceWidth+1) bm.Attrs = map[string]interface{}{"a": "b", "c": int64(1), "d": true} @@ -555,8 +556,8 @@ func TestHandler_Query_Pairs_JSON(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return []interface{}{[]pilosa.Pair{ {ID: 1, Count: 2}, @@ -579,8 +580,8 @@ func TestHandler_Query_Pairs_Protobuf(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return []interface{}{[]pilosa.Pair{ {ID: 1, Count: 2}, @@ -612,15 +613,15 @@ func TestHandler_Query_Err_JSON(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return nil, errors.New("marker") } w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i/query", strings.NewReader(`Bitmap(id=100)`))) - if w.Code != http.StatusInternalServerError { + if w.Code != http.StatusBadRequest { t.Fatalf("unexpected status code: %d", w.Code) } else if body := w.Body.String(); body != `{"error":"marker"}`+"\n" { t.Fatalf("unexpected body: %q", body) @@ -633,8 +634,8 @@ func TestHandler_Query_Err_Protobuf(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return nil, errors.New("marker") } @@ -643,7 +644,7 @@ func TestHandler_Query_Err_Protobuf(t *testing.T) { r := test.MustNewHTTPRequest("POST", "/index/i/query", strings.NewReader(`TopN(frame=x, n=2)`)) r.Header.Set("Accept", "application/x-protobuf") h.ServeHTTP(w, r) - if w.Code != http.StatusInternalServerError { + if w.Code != http.StatusBadRequest { t.Fatalf("unexpected status code: %d", w.Code) } @@ -661,8 +662,8 @@ func TestHandler_Query_MethodNotAllowed(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/index/i/query", nil)) if w.Code != http.StatusMethodNotAllowed { @@ -676,8 +677,8 @@ func TestHandler_Query_ErrParse(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/idx0/query?slices=0,1", strings.NewReader("bad_fn("))) if w.Code != http.StatusBadRequest { @@ -693,7 +694,7 @@ func TestHandler_Index_Delete(t *testing.T) { defer hldr.Close() s := test.NewServer() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder defer s.Close() // Create index. @@ -733,8 +734,8 @@ func TestHandler_DeleteFrame(t *testing.T) { } h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("DELETE", "/index/i0/frame/f1", strings.NewReader(""))) if w.Code != http.StatusOK { @@ -753,8 +754,8 @@ func TestHandler_SetIndexTimeQuantum(t *testing.T) { hldr.MustCreateIndexIfNotExists("i0", pilosa.IndexOptions{}) h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("PATCH", "/index/i0/time-quantum", strings.NewReader(`{"timeQuantum":"ymdh"}`))) if w.Code != http.StatusOK { @@ -777,8 +778,8 @@ func TestHandler_SetFrameTimeQuantum(t *testing.T) { } h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("PATCH", "/index/i0/frame/f1/time-quantum", strings.NewReader(`{"timeQuantum":"ymdh"}`))) if w.Code != http.StatusOK { @@ -796,7 +797,7 @@ func TestHandler_Index_AttrStore_Diff(t *testing.T) { defer hldr.Close() s := test.NewServer() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder defer s.Close() // Set attributes on the index. @@ -845,7 +846,7 @@ func TestHandler_Frame_AttrStore_Diff(t *testing.T) { defer hldr.Close() s := test.NewServer() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder defer s.Close() // Set attributes on the index. @@ -895,7 +896,7 @@ func TestHandler_Frame_AddField(t *testing.T) { defer hldr.Close() s := test.NewServer() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder defer s.Close() t.Run("OK", func(t *testing.T) { @@ -999,7 +1000,7 @@ func TestHandler_Frame_DeleteField(t *testing.T) { defer hldr.Close() s := test.NewServer() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder defer s.Close() t.Run("OK", func(t *testing.T) { @@ -1064,7 +1065,7 @@ func TestHandler_Frame_GetFields(t *testing.T) { defer hldr.Close() s := test.NewServer() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder defer s.Close() t.Run("OK", func(t *testing.T) { @@ -1132,7 +1133,7 @@ func TestHandler_Fragment_BackupRestore(t *testing.T) { defer hldr.Close() s := test.NewServer() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder defer s.Close() // Set bits in the index. @@ -1181,8 +1182,8 @@ func TestHandler_Version(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder w := httptest.NewRecorder() r := test.MustNewHTTPRequest("GET", "/version", nil) @@ -1204,9 +1205,9 @@ func TestHandler_Fragment_Nodes(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(3) - h.Cluster.ReplicaN = 2 + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(3) + h.API.Cluster.ReplicaN = 2 w := httptest.NewRecorder() r := test.MustNewHTTPRequest("GET", "/fragment/nodes?index=X&slice=0", nil) @@ -1233,8 +1234,8 @@ func TestHandler_Expvars(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Cluster = test.NewCluster(1) - h.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder w := httptest.NewRecorder() r := test.MustNewHTTPRequest("GET", "/debug/vars", nil) h.ServeHTTP(w, r) @@ -1279,8 +1280,8 @@ func TestHandler_CreateInputDefinition(t *testing.T) { ] }`) h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i0/input-definition/input1", bytes.NewBuffer(inputBody))) if w.Code != http.StatusOK { @@ -1314,8 +1315,8 @@ func TestHandler_DuplicatePrimaryKey(t *testing.T) { defer hldr.Close() hldr.MustCreateIndexIfNotExists("i0", pilosa.IndexOptions{}) h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) //Ensure throwing error if there's duplicated primaryKey field invalidPrimaryKey := []byte(` @@ -1419,8 +1420,8 @@ func TestHandler_DeleteInputDefinition(t *testing.T) { hldr := test.MustOpenHolder() defer hldr.Close() h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) // Test index not found. w := httptest.NewRecorder() @@ -1467,8 +1468,8 @@ func TestHandler_GetInputDefinition(t *testing.T) { hldr := test.MustOpenHolder() defer hldr.Close() h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) frames := internal.Frame{Name: "f", Meta: &internal.FrameMeta{}} action := internal.InputDefinitionAction{Frame: "f", ValueDestination: "mapping", ValueMap: map[string]uint64{"Green": 1}} @@ -1634,8 +1635,8 @@ func TestHandler_CreateInput(t *testing.T) { "null_value": null }]`) h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) // Return error if index does not exist. w := httptest.NewRecorder() @@ -1749,8 +1750,8 @@ func TestInput_JSON(t *testing.T) { err: "set-timestamp value must be in time format: YYYY-MM-DD, has: 12345"}, } h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) for _, req := range tests { w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i0/input/input1", bytes.NewBuffer([]byte(req.json)))) @@ -1810,8 +1811,8 @@ func TestHandler_DeleteView(t *testing.T) { hldr.Index("i0").Frame("f0").SetTimeQuantum("YMD") h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("DELETE", "/index/i0/frame/f0/view/standard_2017", strings.NewReader(""))) if w.Code != http.StatusOK { @@ -1836,8 +1837,8 @@ func TestHandler_RecalculateCaches(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/recalculate-caches", nil)) @@ -1852,8 +1853,8 @@ func TestHandler_WebUI(t *testing.T) { defer hldr.Close() h := test.NewHandler() - h.Holder = hldr.Holder - h.Cluster = test.NewCluster(1) + h.API.Holder = hldr.Holder + h.API.Cluster = test.NewCluster(1) h.FileSystem = &statik.FileSystem{} w := httptest.NewRecorder() diff --git a/holder_test.go b/holder_test.go index f8766b002..24594025d 100644 --- a/holder_test.go +++ b/holder_test.go @@ -393,7 +393,7 @@ func TestHolderSyncer_SyncHolder(t *testing.T) { defer hldr1.Close() s := test.NewServer() defer s.Close() - s.Handler.Holder = hldr1.Holder + s.Handler.API.Holder = hldr1.Holder s.Handler.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { e := pilosa.NewExecutor(client) e.Holder = hldr1.Holder diff --git a/pilosa.go b/pilosa.go index 8f1ae287e..4b28daccb 100644 --- a/pilosa.go +++ b/pilosa.go @@ -72,6 +72,14 @@ var ( ErrTooManyWrites = errors.New("too many write commands") ErrConfigClusterEnabledHosts = errors.New("providing hosts to a non-disabled cluster is not allowed") + ErrConfigClusterTypeInvalid = errors.New("invalid cluster type") + ErrConfigHostsMissing = errors.New("missing bind address in cluster hosts") + + ErrClusterDoesNotOwnSlice = errors.New("cluster does not own slice") + + ErrNodeIDNotExists = errors.New("node with provided ID does not exist") + ErrNodeNotCoordinator = errors.New("node is not the coordinator") + ErrResizeNotRunning = errors.New("no resize job currently running") ) // Regular expression to validate index and frame names. diff --git a/server.go b/server.go index 8cf2a18ef..21f974ad8 100644 --- a/server.go +++ b/server.go @@ -115,13 +115,14 @@ func NewServer() *Server { Logger: NopLogger, } - s.Handler.Holder = s.Holder - s.diagnostics.server = s + s.Handler.API = NewAPI() + s.Handler.API.Holder = s.Holder return s } // Open opens and initializes the server. func (s *Server) Open() error { + s.Handler.API.Logger = s.Logger // TODO do this in NewServer with functional options s.Logger.Printf("open server") // s.ln can be configured prior to Open() via s.OpenListener(). if s.ln == nil { @@ -164,14 +165,15 @@ func (s *Server) Open() error { s.Cluster.MaxWritesPerRequest = s.MaxWritesPerRequest // Initialize HTTP handler. - s.Handler.Broadcaster = s.Broadcaster - s.Handler.BroadcastHandler = s - s.Handler.StatusHandler = s - s.Handler.Node = node - s.Handler.Cluster = s.Cluster + s.Handler.API.Broadcaster = s.Broadcaster + s.Handler.API.BroadcastHandler = s + s.Handler.API.StatusHandler = s + s.Handler.API.URI = s.URI + s.Handler.API.Cluster = s.Cluster s.Handler.Executor = e s.Cluster.prefect = s.Handler + s.Handler.API.Executor = e // Initialize Holder. s.Holder.Broadcaster = s.Broadcaster diff --git a/server/cluster_test.go b/server/cluster_test.go index 05c971a91..04d53b438 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -62,7 +62,7 @@ func TestMain_SendReceiveMessage(t *testing.T) { m0.Server.Cluster.MemberSet = gossipMemberSet0 m0.Server.Broadcaster = m0.Server m0.Server.Gossiper = gossipMemberSet0 - m0.Server.Handler.Broadcaster = m0.Server.Broadcaster + m0.Server.Handler.API.Broadcaster = m0.Server.Broadcaster m0.Server.Holder.Broadcaster = m0.Server.Broadcaster m0.Server.BroadcastReceiver = gossipMemberSet0 @@ -89,7 +89,7 @@ func TestMain_SendReceiveMessage(t *testing.T) { m1.Server.Cluster.MemberSet = gossipMemberSet1 m1.Server.Broadcaster = m1.Server m1.Server.Gossiper = gossipMemberSet1 - m1.Server.Handler.Broadcaster = m1.Server.Broadcaster + m1.Server.Handler.API.Broadcaster = m1.Server.Broadcaster m1.Server.Holder.Broadcaster = m1.Server.Broadcaster m1.Server.BroadcastReceiver = gossipMemberSet1 @@ -507,9 +507,9 @@ func TestClusterResize_RemoveNode(t *testing.T) { t.Run("ErrorRemoveInvalidNode", func(t *testing.T) { resp := test.MustDo("POST", m0.URL()+fmt.Sprintf("/cluster/resize/remove-node"), `{"id": "invalid-node-id"}`) - expBody := "Node is not a member of the cluster: invalid-node-id" - if resp.StatusCode != http.StatusBadRequest { - t.Fatalf("expected StatusCode %d but got %d", http.StatusBadRequest, resp.StatusCode) + expBody := "removing node: finding node to remove: node with provided ID does not exist" + if resp.StatusCode != http.StatusNotFound { + t.Fatalf("expected StatusCode %d but got %d", http.StatusNotFound, resp.StatusCode) } else if strings.TrimSpace(resp.Body) != expBody { t.Fatalf("expected Body '%s' but got '%s'", expBody, strings.TrimSpace(resp.Body)) } @@ -521,7 +521,7 @@ func TestClusterResize_RemoveNode(t *testing.T) { resp = test.MustDo("POST", m0.URL()+fmt.Sprintf("/cluster/resize/remove-node"), fmt.Sprintf(`{"id": "%s"}`, nodeID)) - expBody := "The coordinator node cannot be removed. First, make a different node the new coordinator." + expBody := "removing node: calling node leave: coordinator cannot be removed; first, make a different node the new coordinator." if resp.StatusCode != http.StatusInternalServerError { t.Fatalf("expected StatusCode %d but got %d", http.StatusInternalServerError, resp.StatusCode) } else if strings.TrimSpace(resp.Body) != expBody { @@ -538,7 +538,7 @@ func TestClusterResize_RemoveNode(t *testing.T) { resp = test.MustDo("POST", m1.URL()+fmt.Sprintf("/cluster/resize/remove-node"), fmt.Sprintf(`{"id": "%s"}`, nodeID)) - expBody := fmt.Sprintf("Node removal requests are only valid on the Coordinator node: %s", coordinatorNodeID) + expBody := fmt.Sprintf("removing node: calling node leave: node removal requests are only valid on the coordinator node: %s", coordinatorNodeID) if resp.StatusCode != http.StatusInternalServerError { t.Fatalf("expected StatusCode %d but got %d", http.StatusInternalServerError, resp.StatusCode) } else if strings.TrimSpace(resp.Body) != expBody { diff --git a/server/server.go b/server/server.go index 66725f4ee..62c7a70a5 100644 --- a/server/server.go +++ b/server/server.go @@ -218,7 +218,7 @@ func (m *Command) SetupServer() error { } c := pilosa.GetHTTPClient(TLSConfig) m.Server.RemoteClient = c - m.Server.Handler.RemoteClient = c + m.Server.Handler.API.RemoteClient = c m.Server.Cluster.RemoteClient = c // Statik file system. diff --git a/server/server_test.go b/server/server_test.go index b91cb4317..da62273a4 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -375,7 +375,11 @@ func TestMain_RecalculateHashes(t *testing.T) { } // Calculate caches on the first node - cluster[0].RecalculateCaches() + err := cluster[0].RecalculateCaches() + if err != nil { + t.Fatalf("recalculating caches: %v", err) + } + target := `{"results":[[{"id":7,"count":99},{"id":1,"count":99},{"id":9,"count":99},{"id":5,"count":99},{"id":4,"count":99},{"id":8,"count":99},{"id":2,"count":99},{"id":6,"count":99},{"id":3,"count":99}]]}` // Run a TopN query on all nodes. The result should be the same as the target. diff --git a/stats_test.go b/stats_test.go index d15d5079b..d6c748c5f 100644 --- a/stats_test.go +++ b/stats_test.go @@ -215,10 +215,10 @@ func TestStatsCount_CreateIndex(t *testing.T) { hldr := test.MustOpenHolder() defer hldr.Close() s := test.NewServer() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder defer s.Close() called := false - s.Handler.Holder.Stats = &MockStats{ + s.Handler.API.Holder.Stats = &MockStats{ mockCount: func(name string, value int64, rate float64) { if name != "createIndex" { t.Errorf("Expected createIndex, Results %s", name) @@ -239,7 +239,7 @@ func TestStatsCount_DeleteIndex(t *testing.T) { defer hldr.Close() s := test.NewServer() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder defer s.Close() // Create index. @@ -247,7 +247,7 @@ func TestStatsCount_DeleteIndex(t *testing.T) { t.Fatal(err) } called := false - s.Handler.Holder.Stats = &MockStats{ + s.Handler.API.Holder.Stats = &MockStats{ mockCount: func(name string, value int64, rate float64) { if name != "deleteIndex" { t.Errorf("Expected deleteIndex, Results %s", name) @@ -268,7 +268,7 @@ func TestStatsCount_CreateFrame(t *testing.T) { defer hldr.Close() s := test.NewServer() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder defer s.Close() // Create index. @@ -276,7 +276,7 @@ func TestStatsCount_CreateFrame(t *testing.T) { t.Fatal(err) } called := false - s.Handler.Holder.Stats = &MockStats{ + s.Handler.API.Holder.Stats = &MockStats{ mockCountWithTags: func(name string, value int64, rate float64, index []string) { if name != "createFrame" { t.Errorf("Expected createFrame, Results %s", name) @@ -300,7 +300,7 @@ func TestStatsCount_DeleteFrame(t *testing.T) { defer hldr.Close() s := test.NewServer() - s.Handler.Holder = hldr.Holder + s.Handler.API.Holder = hldr.Holder defer s.Close() called := false // Create index. @@ -308,7 +308,7 @@ func TestStatsCount_DeleteFrame(t *testing.T) { if _, err := indx.CreateFrameIfNotExists("test", pilosa.FrameOptions{}); err != nil { t.Fatal(err) } - s.Handler.Holder.Stats = &MockStats{ + s.Handler.API.Holder.Stats = &MockStats{ mockCountWithTags: func(name string, value int64, rate float64, index []string) { if name != "deleteFrame" { t.Errorf("Expected deleteFrame, Results %s", name) diff --git a/test/handler.go b/test/handler.go index 0f8c1cefd..cb90cd383 100644 --- a/test/handler.go +++ b/test/handler.go @@ -40,10 +40,12 @@ func NewHandler() *Handler { h := &Handler{ Handler: pilosa.NewHandler(), } - h.Handler.Executor = &h.Executor + h.API = pilosa.NewAPI() + h.Handler.API = h.API + h.Handler.API.Executor = &h.Executor // Handler test messages can no-op. - h.Broadcaster = pilosa.NopBroadcaster + h.API.Broadcaster = pilosa.NopBroadcaster h.SetNormal() @@ -80,14 +82,13 @@ func NewServer() *Server { if err != nil { panic(err) } + s.Handler.API.URI = *uri // Handler test messages can no-op. - s.Handler.Broadcaster = pilosa.NopBroadcaster + s.Handler.API.Broadcaster = pilosa.NopBroadcaster // Create a default cluster on the handler - s.Handler.Cluster = NewCluster(1) - s.Handler.Cluster.Nodes[0].URI = *uri - - s.Handler.Node = s.Handler.Cluster.Nodes[0] + s.Handler.API.Cluster = NewCluster(1) + s.Handler.API.Cluster.Nodes[0].URI = s.HostURI() return s }