Merge branch 'master' into cluster-tests

This commit is contained in:
Matt Jaffee 2018-11-09 12:23:21 -06:00
commit 9c9b1c5afe
No known key found for this signature in database
GPG key ID: 08A3DFFF987B11BF
15 changed files with 3431 additions and 1189 deletions

View file

@ -5,10 +5,12 @@ COPY . /go/src/github.com/pilosa/pilosa/
RUN cd /go/src/github.com/pilosa/pilosa \
&& CGO_ENABLED=0 make install-dep install FLAGS="-a"
FROM scratch
FROM alpine:3.8
LABEL maintainer "dev@pilosa.com"
RUN apk add --no-cache curl jq
COPY --from=builder /go/bin/pilosa /pilosa
COPY LICENSE /LICENSE

171
api.go
View file

@ -690,7 +690,8 @@ func (api *API) FieldAttrDiff(_ context.Context, indexName string, fieldName str
// ImportOptions holds the options for the API.Import method.
type ImportOptions struct {
Clear bool
Clear bool
IgnoreKeyCheck bool
}
// ImportOption is a functional option type for API.Import.
@ -703,8 +704,15 @@ func OptImportOptionsClear(c bool) ImportOption {
}
}
func OptImportOptionsIgnoreKeyCheck(b bool) ImportOption {
return func(o *ImportOptions) error {
o.IgnoreKeyCheck = b
return nil
}
}
// Import bulk imports data into a particular index,field,shard.
func (api *API) Import(_ context.Context, req *ImportRequest, opts ...ImportOption) error {
func (api *API) Import(ctx context.Context, req *ImportRequest, opts ...ImportOption) error {
if err := api.validate(apiImport); err != nil {
return errors.Wrap(err, "validating api method")
}
@ -715,34 +723,71 @@ func (api *API) Import(_ context.Context, req *ImportRequest, opts ...ImportOpti
return errors.Wrap(err, "setting up import options")
}
index := api.holder.Index(req.Index)
if index == nil {
return newNotFoundError(ErrIndexNotFound)
}
field, err := api.indexField(req.Index, req.Field, req.Shard)
index, field, err := api.indexField(req.Index, req.Field, req.Shard)
if err != nil {
return errors.Wrap(err, "getting field")
return errors.Wrap(err, "getting index and field")
}
// Translate row keys.
if field.keys() {
if len(req.RowIDs) != 0 {
return errors.New("row ids cannot be used because field uses string keys")
// Unless explicitly ignoring key validation (meaning keys have been
// translated to ids in a previous step at the coordinator node), then
// check to see if keys need translation.
if !options.IgnoreKeyCheck {
// Translate row keys.
if field.keys() {
if len(req.RowIDs) != 0 {
return errors.New("row ids cannot be used because field uses string keys")
}
if req.RowIDs, err = api.holder.translateFile.TranslateRowsToUint64(index.Name(), field.Name(), req.RowKeys); err != nil {
return errors.Wrap(err, "translating rows")
}
}
if req.RowIDs, err = api.holder.translateFile.TranslateRowsToUint64(index.Name(), field.Name(), req.RowKeys); err != nil {
return errors.Wrap(err, "translating rows")
// Translate column keys.
if index.Keys() {
if len(req.ColumnIDs) != 0 {
return errors.New("column ids cannot be used because index uses string keys")
}
if req.ColumnIDs, err = api.holder.translateFile.TranslateColumnsToUint64(index.Name(), req.ColumnKeys); err != nil {
return errors.Wrap(err, "translating columns")
}
}
// For translated data, map the columnIDs to shards. If
// this node does not own the shard, forward to the node that does.
if index.Keys() || field.keys() {
m := make(map[uint64][]Bit)
for i, colID := range req.ColumnIDs {
shard := colID / ShardWidth
if _, ok := m[shard]; !ok {
m[shard] = make([]Bit, 0)
}
m[shard] = append(m[shard], Bit{
RowID: req.RowIDs[i],
ColumnID: colID,
Timestamp: req.Timestamps[i],
})
}
// Signal to the receiving nodes to ignore checking for key translation.
opts = append(opts, OptImportOptionsIgnoreKeyCheck(true))
var eg errgroup.Group
for shard, bits := range m {
// TODO: if local node owns this shard we don't need to go through the client
shard := shard
bits := bits
eg.Go(func() error {
return api.server.defaultClient.Import(ctx, req.Index, req.Field, shard, bits, opts...)
})
}
return eg.Wait()
}
}
// Translate column keys.
if index.Keys() {
if len(req.ColumnIDs) != 0 {
return errors.New("column ids cannot be used because index uses string keys")
}
if req.ColumnIDs, err = api.holder.translateFile.TranslateColumnsToUint64(index.Name(), req.ColumnKeys); err != nil {
return errors.Wrap(err, "translating columns")
}
// Validate shard ownership.
if err := api.validateShardOwnership(req.Index, req.Shard); err != nil {
return errors.Wrap(err, "validating shard ownership")
}
// Convert timestamps to time.Time.
@ -772,7 +817,7 @@ func (api *API) Import(_ context.Context, req *ImportRequest, opts ...ImportOpti
}
// ImportValue bulk imports values into a particular field.
func (api *API) ImportValue(_ context.Context, req *ImportValueRequest, opts ...ImportOption) error {
func (api *API) ImportValue(ctx context.Context, req *ImportValueRequest, opts ...ImportOption) error {
if err := api.validate(apiImportValue); err != nil {
return errors.Wrap(err, "validating api method")
}
@ -783,26 +828,60 @@ func (api *API) ImportValue(_ context.Context, req *ImportValueRequest, opts ...
return errors.Wrap(err, "setting up import options")
}
index := api.holder.Index(req.Index)
if index == nil {
return newNotFoundError(ErrIndexNotFound)
}
field, err := api.indexField(req.Index, req.Field, req.Shard)
index, field, err := api.indexField(req.Index, req.Field, req.Shard)
if err != nil {
return errors.Wrap(err, "getting field")
return errors.Wrap(err, "getting index and field")
}
// Translate column keys.
if index.Keys() {
if len(req.ColumnIDs) != 0 {
return errors.New("column ids cannot be used because index uses string keys")
}
if req.ColumnIDs, err = api.holder.translateFile.TranslateColumnsToUint64(index.Name(), req.ColumnKeys); err != nil {
return errors.Wrap(err, "translating columns")
// Unless explicitly ignoring key validation (meaning keys have been
// translate to ids in a previous step at the coordinator node), then
// check to see if keys need translation.
if !options.IgnoreKeyCheck {
// Translate column keys.
if index.Keys() {
if len(req.ColumnIDs) != 0 {
return errors.New("column ids cannot be used because index uses string keys")
}
if req.ColumnIDs, err = api.holder.translateFile.TranslateColumnsToUint64(index.Name(), req.ColumnKeys); err != nil {
return errors.Wrap(err, "translating columns")
}
// For translated data, map the columnIDs to shards. If
// this node does not own the shard, forward to the node that does.
m := make(map[uint64][]FieldValue)
for i, colID := range req.ColumnIDs {
shard := colID / ShardWidth
if _, ok := m[shard]; !ok {
m[shard] = make([]FieldValue, 0)
}
m[shard] = append(m[shard], FieldValue{
Value: req.Values[i],
ColumnID: colID,
})
}
// Signal to the receiving nodes to ignore checking for key translation.
opts = append(opts, OptImportOptionsIgnoreKeyCheck(true))
var eg errgroup.Group
for shard, vals := range m {
// TODO: if local node owns this shard we don't need to go through the client
shard := shard
vals := vals
eg.Go(func() error {
return api.server.defaultClient.ImportValue(ctx, req.Index, req.Field, shard, vals, opts...)
})
}
return eg.Wait()
}
}
// Validate shard ownership.
if err := api.validateShardOwnership(req.Index, req.Shard); err != nil {
return errors.Wrap(err, "validating shard ownership")
}
// Import columnIDs into existence field.
if !options.Clear {
if err := importExistenceColumns(index, req.ColumnIDs); err != nil {
@ -863,28 +942,32 @@ func (api *API) LongQueryTime() time.Duration {
return api.cluster.longQueryTime
}
func (api *API) indexField(indexName string, fieldName string, shard uint64) (*Field, error) {
func (api *API) validateShardOwnership(indexName string, shard uint64) error {
// Validate that this handler owns the shard.
if !api.cluster.ownsShard(api.Node().ID, indexName, shard) {
api.server.logger.Printf("node %s does not own shard %d of index %s", api.Node().ID, shard, indexName)
return nil, ErrClusterDoesNotOwnShard
return ErrClusterDoesNotOwnShard
}
return nil
}
func (api *API) indexField(indexName string, fieldName string, shard uint64) (*Index, *Field, error) {
api.server.logger.Printf("importing: %v %v %v", indexName, fieldName, shard)
// Find the Index.
api.server.logger.Printf("importing: %v %v %v", indexName, fieldName, shard)
index := api.holder.Index(indexName)
if index == nil {
api.server.logger.Printf("fragment error: index=%s, field=%s, shard=%d, err=%s", indexName, fieldName, shard, ErrIndexNotFound.Error())
return nil, newNotFoundError(ErrIndexNotFound)
return nil, nil, newNotFoundError(ErrIndexNotFound)
}
// Retrieve field.
field := index.Field(fieldName)
if field == nil {
api.server.logger.Printf("field error: index=%s, field=%s, shard=%d, err=%s", indexName, fieldName, shard, ErrFieldNotFound.Error())
return nil, ErrFieldNotFound
return nil, nil, ErrFieldNotFound
}
return field, nil
return index, field, nil
}
// SetCoordinator makes a new Node the cluster coordinator.

232
api_test.go Normal file
View file

@ -0,0 +1,232 @@
// 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_test
import (
"context"
"fmt"
"reflect"
"testing"
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/server"
"github.com/pilosa/pilosa/test"
)
func TestAPI_Import(t *testing.T) {
c := test.MustRunCluster(t, 2,
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerNodeID("node0"),
pilosa.OptServerClusterHasher(&offsetModHasher{}),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerNodeID("node1"),
pilosa.OptServerClusterHasher(&offsetModHasher{}),
)},
)
defer c.Close()
m0 := c[0]
m1 := c[1]
t.Run("RowIDColumnKey", func(t *testing.T) {
ctx := context.Background()
index := "rick"
field := "f"
_, err := m0.API.CreateIndex(ctx, index, pilosa.IndexOptions{Keys: true})
if err != nil {
t.Fatalf("creating index: %v", err)
}
_, err = m0.API.CreateField(ctx, index, field, pilosa.OptFieldTypeSet(pilosa.DefaultCacheType, 100))
if err != nil {
t.Fatalf("creating field: %v", err)
}
rowID := uint64(1)
timestamp := int64(0)
// Generate some keyed records.
rowIDs := []uint64{}
colKeys := []string{}
timestamps := []int64{}
for i := 1; i <= 10; i++ {
rowIDs = append(rowIDs, rowID)
timestamps = append(timestamps, timestamp)
colKeys = append(colKeys, fmt.Sprintf("col%d", i))
}
// Import data with keys to the coordinator (node0) and verify that it gets
// translated and forwarded to the owner of shard 0 (node1; because of offsetModHasher)
req := &pilosa.ImportRequest{
Index: index,
Field: field,
Shard: 0,
RowIDs: rowIDs,
ColumnKeys: colKeys,
Timestamps: timestamps,
}
if err := m0.API.Import(ctx, req); err != nil {
t.Fatal(err)
}
pql := fmt.Sprintf("Row(%s=%d)", field, rowID)
// Query node0.
if res, err := m0.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
t.Fatal(err)
} else if keys := res.Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, colKeys) {
t.Fatalf("unexpected column keys: %+v", keys)
}
// Query node1.
if res, err := m1.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
t.Fatal(err)
} else if keys := res.Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, colKeys) {
t.Fatalf("unexpected column keys: %+v", keys)
}
})
t.Run("RowKeyColumnID", func(t *testing.T) {
ctx := context.Background()
index := "rkci"
field := "f"
_, err := m0.API.CreateIndex(ctx, index, pilosa.IndexOptions{Keys: false})
if err != nil {
t.Fatalf("creating index: %v", err)
}
_, err = m0.API.CreateField(ctx, index, field, pilosa.OptFieldTypeSet(pilosa.DefaultCacheType, 100), pilosa.OptFieldKeys())
if err != nil {
t.Fatalf("creating field: %v", err)
}
rowKey := "rowkey"
// Generate some keyed records.
rowKeys := []string{rowKey, rowKey, rowKey}
colIDs := []uint64{1, 2, pilosa.ShardWidth + 1}
timestamps := []int64{0, 0, 0}
// Import data with keys to the coordinator (node0) and verify that it gets
// translated and forwarded to the owner of shard 0 (node1; because of offsetModHasher)
req := &pilosa.ImportRequest{
Index: index,
Field: field,
Shard: 0,
RowKeys: rowKeys,
ColumnIDs: colIDs,
Timestamps: timestamps,
}
if err := m0.API.Import(ctx, req); err != nil {
t.Fatal(err)
}
pql := fmt.Sprintf("Row(%s=%s)", field, rowKey)
// Query node0.
if res, err := m0.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
t.Fatal(err)
} else if columns := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, colIDs) {
t.Fatalf("unexpected column ids: %+v", columns)
}
// Query node1.
if res, err := m1.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
t.Fatal(err)
} else if columns := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, colIDs) {
t.Fatalf("unexpected column ids: %+v", columns)
}
})
}
func TestAPI_ImportValue(t *testing.T) {
c := test.MustRunCluster(t, 2,
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerNodeID("node0"),
pilosa.OptServerClusterHasher(&offsetModHasher{}),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerNodeID("node1"),
pilosa.OptServerClusterHasher(&offsetModHasher{}),
)},
)
defer c.Close()
m0 := c[0]
m1 := c[1]
t.Run("ValColumnKey", func(t *testing.T) {
ctx := context.Background()
index := "valck"
field := "f"
_, err := m0.API.CreateIndex(ctx, index, pilosa.IndexOptions{Keys: true})
if err != nil {
t.Fatalf("creating index: %v", err)
}
_, err = m0.API.CreateField(ctx, index, field, pilosa.OptFieldTypeInt(0, 100))
if err != nil {
t.Fatalf("creating field: %v", err)
}
// Generate some keyed records.
values := []int64{}
colKeys := []string{}
for i := 1; i <= 10; i++ {
values = append(values, int64(i))
colKeys = append(colKeys, fmt.Sprintf("col%d", i))
}
// Import data with keys to the coordinator (node0) and verify that it gets
// translated and forwarded to the owner of shard 0 (node1; because of offsetModHasher)
req := &pilosa.ImportValueRequest{
Index: index,
Field: field,
ColumnKeys: colKeys,
Values: values,
}
if err := m0.API.ImportValue(ctx, req); err != nil {
t.Fatal(err)
}
pql := fmt.Sprintf("Range(%s>0)", field)
// Query node0.
if res, err := m0.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
t.Fatal(err)
} else if keys := res.Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, colKeys) {
t.Fatalf("unexpected column keys: %+v", keys)
}
// Query node1.
if res, err := m1.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
t.Fatal(err)
} else if keys := res.Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, colKeys) {
t.Fatalf("unexpected column keys: %+v", keys)
}
})
}
// offsetModHasher represents a simple, mod-based hashing offset by 1.
type offsetModHasher struct{}
func (*offsetModHasher) Hash(key uint64, n int) int {
return int(key+1) % n
}

View file

@ -4,9 +4,9 @@ package pilosa
import "strconv"
const _apiMethod_name = "apiClusterMessageapiCreateFieldapiCreateIndexapiDeleteFieldapiDeleteIndexapiDeleteViewapiExportCSVapiFragmentBlockDataapiFragmentBlocksapiFieldapiFieldAttrDiffapiImportapiImportValueapiIndexapiIndexAttrDiffapiQueryapiRecalculateCachesapiRemoveNodeapiResizeAbortapiSetCoordinatorapiShardNodesapiViews"
const _apiMethod_name = "apiClusterMessageapiCreateFieldapiCreateIndexapiDeleteFieldapiDeleteAvailableShardapiDeleteIndexapiDeleteViewapiExportCSVapiFragmentBlockDataapiFragmentBlocksapiFieldapiFieldAttrDiffapiImportapiImportValueapiIndexapiIndexAttrDiffapiQueryapiRecalculateCachesapiRemoveNodeapiResizeAbortapiSetCoordinatorapiShardNodesapiViews"
var _apiMethod_index = [...]uint16{0, 17, 31, 45, 59, 73, 86, 98, 118, 135, 143, 159, 168, 182, 190, 206, 214, 234, 247, 261, 278, 291, 299}
var _apiMethod_index = [...]uint16{0, 17, 31, 45, 59, 82, 96, 109, 121, 141, 158, 166, 182, 191, 205, 213, 229, 237, 257, 270, 284, 301, 314, 322}
func (i apiMethod) String() string {
if i < 0 || i >= apiMethod(len(_apiMethod_index)-1) {

View file

@ -10,7 +10,6 @@ nav = [
"Time Quantum",
"Attribute",
"Shard",
"View",
]
+++
@ -26,7 +25,7 @@ Pilosa lays out data first in rows, so queries which get all the set bits in one
Please note that Pilosa is most performant when row and column IDs are sequential starting from 0. You can deviate from this to some degree, but setting a bit with column ID 2<sup>63</sup> on a single-node cluster, for example, will not work well due to memory limitations.
![basic data model diagram](/img/docs/data-model.svg)
![basic data model diagram](/img/docs/data-model.png)
*Basic data model diagram*
### Index
@ -89,14 +88,14 @@ This is one major component of Pilosa's ability to combine relationships from mu
Ranked Fields maintain a sorted cache of column counts by Row ID (yielding the top rows by columns with a bit set in each). This cache facilitates the TopN query. The cache size defaults to 50,000 and can be set at Field creation.
![ranked field diagram](/img/docs/field-ranked.svg)
![ranked field diagram](/img/docs/field-ranked.png)
*Ranked field diagram*
#### LRU
The LRU cache maintains the most recently accessed Rows.
![lru field diagram](/img/docs/field-lru.svg)
![lru field diagram](/img/docs/field-lru.png)
*LRU field diagram*
### Time Quantum
@ -161,7 +160,7 @@ Set(2, B=1)
Set(3, B=6)
```
![BSI field diagram](/img/docs/field-bsi.svg)
![BSI field diagram](/img/docs/field-bsi.png)
*BSI field diagram*
Check out this [blog post](/blog/range-encoded-bitmaps/) for some more details about BSI in Pilosa.
@ -186,7 +185,7 @@ Set(3, A=8, 2017-05-18T00:00)
Set(3, A=8, 2017-05-19T00:00)
```
![time quantum field diagram](/img/docs/field-time-quantum.svg)
![time quantum field diagram](/img/docs/field-time-quantum.png)
*Time quantum fueld diagram*
#### Mutex

View file

@ -70,4 +70,4 @@ nav = []
<strong id="topn">[TopN](../query-language/#topn):</strong> A [PQL](#pql) query that returns a list of rows, sorted by the count of [columns](#column) set in the [row](#row), within a specified [field](#field).
<strong id="view">[View](../data-model/#view):</strong> Views separate the different data layouts within a [Field](#field). The primary view is standard, which represents the typical [row](#row)/[column](#column) data. Time based field views are automatically generated for each [time quantum](#time-quantum). Views are internally managed by Pilosa, and never exposed directly via the API. This simplifies the functional interface by separating it from the physical data representation.
<strong id="view">View:</strong> Views separate the different data layouts within a [Field](#field). The primary view is standard, which represents the typical [row](#row)/[column](#column) data. Time based field views are automatically generated for each [time quantum](#time-quantum). Views are internally managed by Pilosa, and never exposed directly via the API. This simplifies the functional interface by separating it from the physical data representation.

View file

@ -75,7 +75,7 @@ type (
// 0 if a == b
// > 0 if a > b
//
Cmp func(a, b uint64) int
Cmp func(a, b uint64) int64
d struct { // data page
c int

View file

@ -23,8 +23,8 @@ import (
"github.com/pilosa/pilosa/roaring"
)
func cmp(a, b uint64) int {
return int(a - b)
func cmp(a, b uint64) int64 {
return int64(a - b)
}
type bTreeContainers struct {

View file

@ -84,9 +84,9 @@ func (c *InternalClient) maxShardByIndex(ctx context.Context) (map[string]uint64
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, errors.Wrap(err, "executing request")
return nil, err
}
defer resp.Body.Close()
@ -115,16 +115,14 @@ func (c *InternalClient) Schema(ctx context.Context) ([]*pilosa.IndexInfo, error
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, errors.Wrap(err, "executing request")
return nil, err
}
defer resp.Body.Close()
var rsp getSchemaResponse
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("http: status=%d", resp.StatusCode)
} else if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil {
if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil {
return nil, fmt.Errorf("json decode: %s", err)
}
return rsp.Indexes, nil
@ -152,27 +150,14 @@ func (c *InternalClient) CreateIndex(ctx context.Context, index string, opt pilo
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
// Execute request against the host.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return errors.Wrap(err, "executing request")
}
defer resp.Body.Close()
// Read body.
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return errors.Wrap(err, "reading")
}
// Handle response based on status code.
switch resp.StatusCode {
case http.StatusOK:
return nil // ok
case http.StatusConflict:
return pilosa.ErrIndexExists
default:
return errors.New(string(body))
if resp.StatusCode == http.StatusConflict {
return pilosa.ErrIndexExists
}
return err
}
return nil
}
// FragmentNodes returns a list of nodes that own a shard.
@ -191,19 +176,16 @@ func (c *InternalClient) FragmentNodes(ctx context.Context, index string, shard
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, errors.Wrap(err, "executing request")
return nil, err
}
defer resp.Body.Close()
var a []*pilosa.Node
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("http: status=%d", resp.StatusCode)
} else if err := json.NewDecoder(resp.Body).Decode(&a); err != nil {
if err := json.NewDecoder(resp.Body).Decode(&a); err != nil {
return nil, fmt.Errorf("json decode: %s", err)
}
return a, nil
}
@ -222,19 +204,16 @@ func (c *InternalClient) Nodes(ctx context.Context) ([]*pilosa.Node, error) {
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, errors.Wrap(err, "executing request")
return nil, err
}
defer resp.Body.Close()
var a []*pilosa.Node
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("http: status=%d", resp.StatusCode)
} else if err := json.NewDecoder(resp.Body).Decode(&a); err != nil {
if err := json.NewDecoder(resp.Body).Decode(&a); err != nil {
return nil, fmt.Errorf("json decode: %s", err)
}
return a, nil
}
@ -269,9 +248,9 @@ func (c *InternalClient) QueryNode(ctx context.Context, uri *pilosa.URI, index s
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
// Execute request against the host.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, errors.Wrap(err, "executing request")
return nil, err
}
defer resp.Body.Close()
@ -437,6 +416,9 @@ func (c *InternalClient) importNode(ctx context.Context, node *pilosa.Node, inde
if opts.Clear {
vals.Set("clear", "true")
}
if opts.IgnoreKeyCheck {
vals.Set("ignoreKeyCheck", "true")
}
url := fmt.Sprintf("%s?%s", u.String(), vals.Encode())
req, err := http.NewRequest("POST", url, bytes.NewReader(buf))
@ -449,9 +431,9 @@ func (c *InternalClient) importNode(ctx context.Context, node *pilosa.Node, inde
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
// Execute request against the host.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return errors.Wrap(err, "executing request")
return err
}
defer resp.Body.Close()
@ -459,8 +441,6 @@ func (c *InternalClient) importNode(ctx context.Context, node *pilosa.Node, inde
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return errors.Wrap(err, "reading")
} else if resp.StatusCode != http.StatusOK {
return errors.New(string(body))
}
var isresp pilosa.ImportResponse
@ -605,17 +585,12 @@ func (c *InternalClient) ImportRoaring(ctx context.Context, uri *pilosa.URI, ind
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
// Execute request against the host.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return errors.Wrap(err, "executing request")
return err
}
defer resp.Body.Close()
// Validate status code.
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("invalid status: %d", resp.StatusCode)
}
dec := json.NewDecoder(resp.Body)
rbody := &pilosa.ImportResponse{}
dec.Decode(rbody)
@ -674,17 +649,12 @@ func (c *InternalClient) exportNodeCSV(ctx context.Context, node *pilosa.Node, i
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
// Execute request against the host.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return errors.Wrap(err, "executing request")
return err
}
defer resp.Body.Close()
// Validate status code.
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("invalid status: %d", resp.StatusCode)
}
// Copy body to writer.
if _, err := io.Copy(w, resp.Body); err != nil {
return errors.Wrap(err, "copying")
@ -717,18 +687,12 @@ func (c *InternalClient) backupShardNode(ctx context.Context, index, field strin
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
// Execute request.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, errors.Wrap(err, "executing request")
}
// Return error if status is not OK.
if resp.StatusCode == http.StatusNotFound {
resp.Body.Close()
return nil, pilosa.ErrFragmentNotFound
} else if resp.StatusCode != http.StatusOK {
resp.Body.Close()
return nil, fmt.Errorf("unexpected backup status code: host=%s, code=%d", node.URI, resp.StatusCode)
if resp.StatusCode == http.StatusNotFound {
return nil, pilosa.ErrFragmentNotFound
}
return nil, err
}
return resp.Body, nil
@ -780,27 +744,15 @@ func (c *InternalClient) CreateFieldWithOptions(ctx context.Context, index, fiel
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
// Execute request against the host.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return errors.Wrap(err, "executing request")
}
defer resp.Body.Close()
// Read body.
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return errors.Wrap(err, "reading")
if resp.StatusCode == http.StatusConflict {
return pilosa.ErrFieldExists
}
return err
}
// Handle response based on status code.
switch resp.StatusCode {
case http.StatusOK:
return nil // ok
case http.StatusConflict:
return pilosa.ErrFieldExists
default:
return errors.New(string(body))
}
return nil
}
// FragmentBlocks returns a list of block checksums for a fragment on a host.
@ -827,21 +779,16 @@ func (c *InternalClient) FragmentBlocks(ctx context.Context, uri *pilosa.URI, in
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, errors.Wrap(err, "executing request")
// Return the appropriate error.
if resp.StatusCode == http.StatusNotFound {
return nil, pilosa.ErrFragmentNotFound
}
return nil, err
}
defer resp.Body.Close()
// Return error if status is not OK.
switch resp.StatusCode {
case http.StatusOK: // ok
case http.StatusNotFound:
return nil, pilosa.ErrFragmentNotFound
default:
return nil, fmt.Errorf("unexpected status: code=%d", resp.StatusCode)
}
// Decode response object.
var rsp getFragmentBlocksResponse
if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil {
@ -876,21 +823,15 @@ func (c *InternalClient) BlockData(ctx context.Context, uri *pilosa.URI, index,
req.Header.Set("Accept", "application/protobuf")
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, nil, errors.Wrap(err, "executing request")
if resp.StatusCode == http.StatusNotFound {
return nil, nil, nil
}
return nil, nil, err
}
defer resp.Body.Close()
// Return error if status is not OK.
switch resp.StatusCode {
case http.StatusOK: // fallthrough
case http.StatusNotFound:
return nil, nil, nil
default:
return nil, nil, fmt.Errorf("unexpected status: code=%d", resp.StatusCode)
}
// Decode response object.
var rsp pilosa.BlockDataResponse
if body, err := ioutil.ReadAll(resp.Body); err != nil {
@ -924,19 +865,12 @@ func (c *InternalClient) ColumnAttrDiff(ctx context.Context, uri *pilosa.URI, in
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, errors.Wrap(err, "executing request")
return nil, err
}
defer resp.Body.Close()
// Return error if status is not OK.
switch resp.StatusCode {
case http.StatusOK: // ok
default:
return nil, fmt.Errorf("unexpected status: code=%d", resp.StatusCode)
}
// Decode response object.
var rsp postIndexAttrDiffResponse
if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil {
@ -968,21 +902,15 @@ func (c *InternalClient) RowAttrDiff(ctx context.Context, uri *pilosa.URI, index
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.httpClient.Do(req.WithContext(ctx))
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, errors.Wrap(err, "executing request")
if resp.StatusCode == http.StatusNotFound {
return nil, pilosa.ErrFieldNotFound
}
return nil, err
}
defer resp.Body.Close()
// Return error if status is not OK.
switch resp.StatusCode {
case http.StatusOK: // ok
case http.StatusNotFound:
return nil, pilosa.ErrFieldNotFound
default:
return nil, fmt.Errorf("unexpected status: code=%d", resp.StatusCode)
}
// Decode response object.
var rsp postFieldAttrDiffResponse
if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil {
@ -1003,26 +931,33 @@ func (c *InternalClient) SendMessage(ctx context.Context, uri *pilosa.URI, msg [
req.Header.Set("Accept", "application/json")
// Execute request.
resp, err := c.httpClient.Do(req.WithContext(ctx))
_, err = c.executeRequest(req.WithContext(ctx))
return err
}
// executeRequest executes the given request and checks the Response
func (c *InternalClient) executeRequest(req *http.Request) (*http.Response, error) {
resp, err := c.httpClient.Do(req)
if err != nil {
return fmt.Errorf("executing http request: %v", err)
return nil, errors.Wrap(err, "executing request")
}
defer resp.Body.Close()
// Read body.
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return fmt.Errorf("reading response body: %v", err)
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
defer resp.Body.Close()
buf, err := ioutil.ReadAll(resp.Body)
if err != nil {
return resp, errors.Wrapf(err, "bad status '%s' and err reading body", resp.Status)
}
var msg string
// try to decode a JSON response
var sr successResponse
if err = json.Unmarshal(buf, &sr); err == nil {
msg = sr.Error.Error()
} else {
msg = string(buf)
}
return resp, errors.Errorf("server error %s: '%s'", resp.Status, msg)
}
// Return error if status is not OK.
switch resp.StatusCode {
case http.StatusOK: // ok
default:
return fmt.Errorf("unexpected response status code: %d: %s", resp.StatusCode, body)
}
return nil
return resp, nil
}
// Bits is a slice of Bit.

View file

@ -181,7 +181,7 @@ func (h *Handler) populateValidators() {
h.validators["DeleteIndex"] = queryValidationSpecRequired()
h.validators["PostField"] = queryValidationSpecRequired()
h.validators["DeleteField"] = queryValidationSpecRequired()
h.validators["PostImport"] = queryValidationSpecRequired().Optional("clear")
h.validators["PostImport"] = queryValidationSpecRequired().Optional("clear", "ignoreKeyCheck")
h.validators["PostImportRoaring"] = queryValidationSpecRequired().Optional("remote", "clear")
h.validators["PostQuery"] = queryValidationSpecRequired().Optional("shards", "columnAttrs", "excludeRowAttrs", "excludeColumns")
h.validators["GetInfo"] = queryValidationSpecRequired()
@ -260,20 +260,10 @@ func newRouter(handler *Handler) *mux.Router {
router.HandleFunc("/internal/shards/max", handler.handleGetShardsMax).Methods("GET").Name("GetShardsMax") // TODO: deprecate, but it's being used by the client
router.HandleFunc("/internal/translate/data", handler.handleGetTranslateData).Methods("GET").Name("GetTranslateData")
// TODO: Apply MethodNotAllowed statuses to all endpoints.
// Ideally this would be automatic, as described in this (wontfix) ticket:
// https://github.com/gorilla/mux/issues/6
// For now we just do it for the most commonly used handler, /query
router.HandleFunc("/index/{index}/query", handler.methodNotAllowedHandler).Methods("GET")
router.Use(handler.queryArgValidator)
return router
}
func (h *Handler) methodNotAllowedHandler(w http.ResponseWriter, _ *http.Request) {
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
}
// ServeHTTP handles an HTTP request.
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
defer func() {
@ -997,6 +987,12 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) {
// If the clear flag is true, treat the import as clear bits.
q := r.URL.Query()
doClear := q.Get("clear") == "true"
doIgnoreKeyCheck := q.Get("ignoreKeyCheck") == "true"
opts := []pilosa.ImportOption{
pilosa.OptImportOptionsClear(doClear),
pilosa.OptImportOptionsIgnoreKeyCheck(doIgnoreKeyCheck),
}
// Get index and field type to determine how to handle the
// import data.
@ -1030,7 +1026,7 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) {
return
}
if err := h.api.ImportValue(r.Context(), req, pilosa.OptImportOptionsClear(doClear)); err != nil {
if err := h.api.ImportValue(r.Context(), req, opts...); err != nil {
switch errors.Cause(err) {
case pilosa.ErrClusterDoesNotOwnShard:
http.Error(w, err.Error(), http.StatusPreconditionFailed)
@ -1048,7 +1044,7 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) {
return
}
if err := h.api.Import(r.Context(), req, pilosa.OptImportOptionsClear(doClear)); err != nil {
if err := h.api.Import(r.Context(), req, opts...); err != nil {
switch errors.Cause(err) {
case pilosa.ErrClusterDoesNotOwnShard:
http.Error(w, err.Error(), http.StatusPreconditionFailed)

File diff suppressed because it is too large Load diff

File diff suppressed because it is too large Load diff

View file

@ -75,19 +75,20 @@ func TestParser_Parse(t *testing.T) {
// Parse with only arguments.
t.Run("ArgumentsOnly", func(t *testing.T) {
q, err := pql.ParseString(`MyCall( key= value, foo="bar", age = 12 , bool0=true, bool1=false, x=null )`)
q, err := pql.ParseString(`MyCall( key= value, foo='bar', age = 12 , bool0=true, bool1=false, x=null, escape="\" \\escape\n\\\\" )`)
if err != nil {
t.Fatal(err)
} else if !reflect.DeepEqual(q.Calls[0],
&pql.Call{
Name: "MyCall",
Args: map[string]interface{}{
"key": "value",
"foo": "bar",
"age": int64(12),
"bool0": true,
"bool1": false,
"x": nil,
"key": "value",
"foo": "bar",
"age": int64(12),
"bool0": true,
"bool1": false,
"x": nil,
"escape": "\" \\escape\n\\\\",
},
},
) {

View file

@ -45,19 +45,18 @@ item <- ( 'null' &(comma / sp close) { p.addVal(nil) }
/ < '-'? [0-9]+ ('.'[0-9]*)? > { p.addNumVal(buffer[begin:end]) }
/ < '-'? '.'[0-9]+ > { p.addNumVal(buffer[begin:end]) }
/ < ([[A-Z]] / [0-9] / '-' / '_' / ':')+ > { p.addVal(buffer[begin:end]) }
/ '"' < doublequotedstring > '"' { p.addVal(buffer[begin:end]) }
/ < '"' doublequotedstring '"' > { s, _ := strconv.Unquote(buffer[begin:end]); p.addVal(s) }
/ '\'' < singlequotedstring > '\'' { p.addVal(buffer[begin:end]) }
)
doublequotedstring <- ( [^"\\\n] / '\\n' / '\\\"' / '\\\'' / '\\\\' )*
singlequotedstring <- ( [^'\\\n] / '\\n' / '\\\"' / '\\\'' / '\\\\' )*
doublequotedstring <- ( '\\"' / '\\\\' / [^"] )*
singlequotedstring <- ( '\\\'' / '\\\\' / [^'] )*
fieldExpr <- [[A-Z]] ( [[A-Z]] / [0-9] / '_' / '-' )*
field <- <fieldExpr / reserved> { p.addField(buffer[begin:end]) }
reserved <- ('_row' / '_col' / '_start' / '_end' / '_timestamp' / '_field')
posfield <- <fieldExpr> { p.addPosStr("_field", buffer[begin:end]) }
uint <- [1-9] [0-9]* / '0'
uintrow <- <uint>{p.addPosNum("_row", buffer[begin:end])}
col <- ( <uint> {p.addPosNum("_col", buffer[begin:end])}
/ '\'' <singlequotedstring> '\'' {p.addPosStr("_col", buffer[begin:end])}
/ '"' <doublequotedstring> '"' {p.addPosStr("_col", buffer[begin:end])}

File diff suppressed because it is too large Load diff