mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 10:54:59 +00:00
Merge pull request #1190 from jaffee/1151-http-refactor
1151 http refactor
This commit is contained in:
commit
dd5b38f4ef
18 changed files with 1360 additions and 927 deletions
896
api.go
Normal file
896
api.go
Normal file
|
|
@ -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()
|
||||
}
|
||||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
||||
|
|
|
|||
|
|
@ -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("")))
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
||||
|
|
|
|||
1009
handler.go
1009
handler.go
File diff suppressed because it is too large
Load diff
169
handler_test.go
169
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()
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
16
server.go
16
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
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue