Merge branch 'master' into shift-op

This commit is contained in:
Travis Turner 2019-01-24 13:56:28 -06:00
commit 2467d88ddc
No known key found for this signature in database
GPG key ID: 7F08008DFD9314C9
35 changed files with 1569 additions and 1206 deletions

View file

@ -139,9 +139,9 @@ gometalinter: require-gometalinter
--enable=ineffassign \
--enable=interfacer \
--enable=maligned \
--enable=megacheck \
--enable=misspell \
--enable=nakedret \
--enable=staticcheck \
--enable=unconvert \
--enable=unparam \
--enable=vet \

View file

@ -609,8 +609,8 @@ func (c *cluster) addNodeBasicSorted(node *Node) bool {
// Nodes returns a copy of the slice of nodes in the cluster. Safe for
// concurrent use, result may be modified.
func (c *cluster) Nodes() []*Node {
c.mu.Lock()
defer c.mu.Unlock()
c.mu.RLock()
defer c.mu.RUnlock()
ret := make([]*Node, len(c.nodes))
copy(ret, c.nodes)
return ret
@ -1792,7 +1792,7 @@ func (c *cluster) nodeLeave(nodeID string) error {
}
if c.state != ClusterStateNormal && c.state != ClusterStateDegraded {
return fmt.Errorf("Cluster must be in state %s to remove a node. Current state: %s",
return fmt.Errorf("cluster must be '%s' to remove a node but is '%s'",
ClusterStateNormal, c.state)
}
@ -1803,7 +1803,7 @@ func (c *cluster) nodeLeave(nodeID string) error {
// Prevent removing the coordinator node (this node).
if nodeID == c.Node.ID {
return fmt.Errorf("coordinator 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

View file

@ -83,7 +83,7 @@ func TestFragCombos(t *testing.T) {
// newIndexWithTempPath returns a new instance of Index.
func newIndexWithTempPath(name string) *Index {
path, err := ioutil.TempDir("", "pilosa-index-")
path, err := ioutil.TempDir(*TempDir, "pilosa-index-")
if err != nil {
panic(err)
}

View file

@ -151,11 +151,11 @@ func (cmd *ImportCommand) Run(ctx context.Context) error {
func (cmd *ImportCommand) ensureSchema(ctx context.Context) error {
err := cmd.client.EnsureIndex(ctx, cmd.Index, cmd.IndexOptions)
if err != nil {
return fmt.Errorf("Error Creating Index: %s", err)
return errors.Wrap(err, "creating index")
}
err = cmd.client.EnsureFieldWithOptions(ctx, cmd.Index, cmd.Field, cmd.FieldOptions)
if err != nil {
return fmt.Errorf("Error Creating Field: %s", err)
return errors.Wrap(err, "creating field")
}
return nil
}

View file

@ -137,11 +137,11 @@ func (d *diagnosticsCollector) compareVersion(value string) error {
localVersion := versionSegments(d.version)
if localVersion[0] < currentVersion[0] { //Major
return fmt.Errorf("Warning: You are running Pilosa %s. A newer version (%s) is available: https://github.com/pilosa/pilosa/releases", d.version, value)
return fmt.Errorf("you are running Pilosa %s, a newer version (%s) is available: https://github.com/pilosa/pilosa/releases", d.version, value)
} else if localVersion[1] < currentVersion[1] && localVersion[0] == currentVersion[0] { // Minor
return fmt.Errorf("Warning: You are running Pilosa %s. The latest Minor release is %s: https://github.com/pilosa/pilosa/releases", d.version, value)
return fmt.Errorf("you are running Pilosa %s, the latest minor release is %s: https://github.com/pilosa/pilosa/releases", d.version, value)
} else if localVersion[2] < currentVersion[2] && localVersion[0] == currentVersion[0] && localVersion[1] == currentVersion[1] { // Patch
return fmt.Errorf("There is a new patch release of Pilosa available: %s: https://github.com/pilosa/pilosa/releases", value)
return fmt.Errorf("there is a new patch release of Pilosa available: %s: https://github.com/pilosa/pilosa/releases", value)
}
return nil

View file

@ -82,19 +82,19 @@ func TestDiagnosticsVersion_Compare(t *testing.T) {
d.SetVersion(version)
err := d.compareVersion("v1.7.0")
if !strings.Contains(err.Error(), "A newer version") {
if !strings.Contains(err.Error(), "a newer version") {
t.Fatalf("Expected a newer version is available, actual error: %s", err)
}
err = d.compareVersion("1.7.0")
if !strings.Contains(err.Error(), "A newer version") {
if !strings.Contains(err.Error(), "a newer version") {
t.Fatalf("Expected a newer version is available, actual error: %s", err)
}
err = d.compareVersion("0.7.0")
if !strings.Contains(err.Error(), "The latest Minor release is") {
if !strings.Contains(err.Error(), "the latest minor release is") {
t.Fatalf("Expected Minor Version Missmatch, actual error: %s", err)
}
err = d.compareVersion("0.1.2")
if !strings.Contains(err.Error(), "There is a new patch release of Pilosa") {
if !strings.Contains(err.Error(), "there is a new patch release of Pilosa") {
t.Fatalf("Expected Patch Version Missmatch, actual error: %s", err)
}
err = d.compareVersion("0.1.1")

View file

@ -366,7 +366,7 @@ A three node cluster running on different hosts could be minimally configured as
[gossip]
port = 12000
seed = "node0.pilosa.com:12000"
seeds = ["node0.pilosa.com:12000"]
[cluster]
replicas = 1
@ -379,7 +379,7 @@ A three node cluster running on different hosts could be minimally configured as
[gossip]
port = 12000
seed = "node0.pilosa.com:12000"
seeds = ["node0.pilosa.com:12000"]
[cluster]
replicas = 1
@ -392,7 +392,7 @@ A three node cluster running on different hosts could be minimally configured as
[gossip]
port = 12000
seed = "node0.pilosa.com:12000"
seeds = ["node0.pilosa.com:12000"]
[cluster]
replicas = 1
@ -410,7 +410,7 @@ The same cluster which uses HTTPS instead of HTTP can be configured as follows.
[gossip]
port = 12000
seed = "node0.pilosa.com:12000"
seeds = ["node0.pilosa.com:12000"]
key = "/home/pilosa/private/gossip.key32"
[cluster]
@ -428,7 +428,7 @@ The same cluster which uses HTTPS instead of HTTP can be configured as follows.
[gossip]
port = 12000
seed = "node0.pilosa.com:12000"
seeds = ["node0.pilosa.com:12000"]
key = "/home/pilosa/private/gossip.key32"
[cluster]
@ -446,7 +446,7 @@ The same cluster which uses HTTPS instead of HTTP can be configured as follows.
[gossip]
port = 12000
seed = "node0.pilosa.com:12000"
seeds = ["node0.pilosa.com:12000"]
key = "/home/pilosa/private/gossip.key32"
[cluster]
@ -468,7 +468,7 @@ You can run a cluster on the same host using the configuration above with a few
[gossip]
port = 12000
seed = "localhost:12000"
seeds = ["localhost:12000"]
key = "/home/pilosa/private/gossip.key32"
[cluster]
@ -486,7 +486,7 @@ You can run a cluster on the same host using the configuration above with a few
[gossip]
port = 12001
seed = "localhost:12000"
seeds = ["localhost:12000"]
key = "/home/pilosa/private/gossip.key32"
[cluster]
@ -504,7 +504,7 @@ You can run a cluster on the same host using the configuration above with a few
[gossip]
port = 12002
seed = "localhost:12000"
seeds = ["localhost:12000"]
key = "/home/pilosa/private/gossip.key32"
[cluster]

View file

@ -786,7 +786,7 @@ Options(Row(f1=10), shards=[0, 2])
**Spec:**
```
Rows(field=<STRING>, previous=<UINT|STRING>, limit=<UINT>, column=<UINT|STRING>)
Rows(<FIELD>, previous=<UINT|STRING>, limit=<UINT>, column=<UINT|STRING>)
```
**Description:**
@ -810,7 +810,7 @@ result of the previous request will start from the next available row.
Without keys:
```request
Rows(field=blah)
Rows(blah)
```
```response
{"rows":[1,9,39]}
@ -818,7 +818,7 @@ Rows(field=blah)
With keys:
```request
Rows(field=blahk)
Rows(blahk)
```
```response
{"rows":null,"keys":["haha","zaaa","traa"]}
@ -859,7 +859,7 @@ specify the field and row for each row that was intersected to get that result.
A single `Rows` query.
```request
GroupBy(Rows(field=blah))
GroupBy(Rows(blah))
```
```response
[{"group":[{"field":"blah","rowID":1}],"count":1},
@ -869,7 +869,7 @@ GroupBy(Rows(field=blah))
With two `Rows` queries - one with IDs and one with keys.
```request
GroupBy(Rows(field=blah), Rows(field=blahk), limit=7)
GroupBy(Rows(blah), Rows(blahk), limit=7)
```
```response
[{"group":[{"field":"blah","rowID":1},{"field":"blahk","rowKey":"haha"}],"count":1},
@ -883,7 +883,7 @@ GroupBy(Rows(field=blah), Rows(field=blahk), limit=7)
Getting the rest of the results from the previous example (paging).
```request
GroupBy(Rows(field=blah, previous=39), Rows(field=blahk, previous="haha"), limit=7)
GroupBy(Rows(blah, previous=39), Rows(blahk, previous="haha"), limit=7)
```
```response

View file

@ -18,7 +18,6 @@ import (
"context"
"encoding/json"
"fmt"
"log"
"sort"
"time"
@ -808,7 +807,7 @@ func (e *executor) executeTopNShard(ctx context.Context, index string, c *pql.Ca
return nil, nil
}
if minThreshold <= 0 {
if minThreshold == 0 {
minThreshold = defaultMinThreshold
}
@ -917,6 +916,12 @@ func (e *executor) executeGroupBy(ctx context.Context, index string, c *pql.Call
// TODO support TopN in here would be really cool - and pretty easy I think.
childRows := make([]RowIDs, len(c.Children))
for i, child := range c.Children {
// Check "field" first for backwards compatibility, then set _field.
// TODO: remove at Pilosa 2.0
if fieldName, ok := child.Args["field"].(string); ok {
child.Args["_field"] = fieldName
}
if child.Name != "Rows" {
return nil, errors.Errorf("'%s' is not a valid child query for GroupBy, must be 'Rows'", child.Name)
}
@ -1091,6 +1096,17 @@ func (e *executor) executeGroupByShard(ctx context.Context, index string, c *pql
}
func (e *executor) executeRows(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *execOptions) (RowIDs, error) {
// Fetch field name from argument.
// Check "field" first for backwards compatibility.
// TODO: remove at Pilosa 2.0
var fieldName string
var ok bool
if fieldName, ok = c.Args["field"].(string); ok {
c.Args["_field"] = fieldName
}
if fieldName, ok = c.Args["_field"].(string); !ok {
return nil, errors.New("Rows() field required")
}
if columnID, ok, err := c.UintArg("column"); err != nil {
return nil, errors.Wrap(err, "getting column")
} else if ok {
@ -1099,7 +1115,7 @@ func (e *executor) executeRows(ctx context.Context, index string, c *pql.Call, s
// Execute calls in bulk on each remote node and merge.
mapFn := func(shard uint64) (interface{}, error) {
return e.executeRowsShard(ctx, index, c, shard)
return e.executeRowsShard(ctx, index, fieldName, c, shard)
}
// Determine limit so we can use it when reducing.
@ -1124,22 +1140,25 @@ func (e *executor) executeRows(ctx context.Context, index string, c *pql.Call, s
return results, nil
}
func (e *executor) executeRowsShard(_ context.Context, index string, c *pql.Call, shard uint64) (RowIDs, error) {
func (e *executor) executeRowsShard(_ context.Context, index string, fieldName string, c *pql.Call, shard uint64) (RowIDs, error) {
// Fetch index.
idx := e.Holder.Index(index)
if idx == nil {
return nil, ErrIndexNotFound
}
// Fetch field name from argument.
fieldName, ok := c.Args["field"].(string)
if !ok {
return nil, errors.New("Rows() argument required: field")
}
// Fetch field.
f := e.Holder.Field(index, fieldName)
if f == nil {
return nil, ErrFieldNotFound
}
// Rows query does not currently support a `time` field that has
// `noStandardView: true`.
// TODO https://github.com/pilosa/pilosa/issues/1783
if f.Type() == FieldTypeTime && f.options.NoStandardView {
return nil, errors.New("Rows() query on time field with no standard view is not currently supported")
}
frag := e.Holder.fragment(index, fieldName, viewStandard, shard)
if frag == nil {
return make(RowIDs, 0), nil
@ -1176,7 +1195,7 @@ func (e *executor) executeRowShard(ctx context.Context, index string, c *pql.Cal
defer span.Finish()
if c.Name == "Range" {
log.Print("DEPRECATED: Range() is deprecated, please use Row() instead.")
e.Holder.Logger.Printf("DEPRECATED: Range() is deprecated, please use Row() instead.")
}
// Handle bsiGroup ranges differently.
@ -1586,14 +1605,14 @@ func (e *executor) executeClearBit(ctx context.Context, index string, c *pql.Cal
if err != nil {
return false, fmt.Errorf("reading Clear() row: %v", err)
} else if !ok {
return false, fmt.Errorf("Clear() row argument '%v' required", rowLabel)
return false, fmt.Errorf("row=<row> argument required to Clear() call")
}
colID, ok, err := c.UintArg("_" + columnLabel)
if err != nil {
return false, fmt.Errorf("reading Clear() column: %v", err)
} else if !ok {
return false, fmt.Errorf("Clear() col argument '%v' required", columnLabel)
return false, fmt.Errorf("column argument to Clear(<COLUMN>, <FIELD>=<ROW>) required")
}
return e.executeClearBitField(ctx, index, c, f, colID, rowID, opt)
@ -1712,19 +1731,19 @@ func (e *executor) executeClearRowShard(ctx context.Context, index string, c *pq
return changed, nil
}
// executeSetRow executes a SetRow() call.
// executeSetRow executes a Store() call.
func (e *executor) executeSetRow(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *execOptions) (bool, error) {
// Ensure the field type supports SetRow().
// Ensure the field type supports Store().
fieldName, err := c.FieldArg()
if err != nil {
return false, errors.New("SetRow() argument required: field")
return false, errors.New("field required for Store()")
}
field := e.Holder.Field(index, fieldName)
if field == nil {
return false, ErrFieldNotFound
}
if field.Type() != FieldTypeSet {
return false, fmt.Errorf("SetRow() is not supported on %s field types", field.Type())
return false, fmt.Errorf("can't Store() on a %s field", field.Type())
}
// Execute calls in bulk on each remote node and merge.
@ -1749,15 +1768,15 @@ func (e *executor) executeSetRow(ctx context.Context, index string, c *pql.Call,
func (e *executor) executeSetRowShard(ctx context.Context, index string, c *pql.Call, shard uint64) (bool, error) {
fieldName, err := c.FieldArg()
if err != nil {
return false, errors.New("SetRow() argument required: field")
return false, errors.New("Store() argument required: field")
}
// Read fields using labels.
rowID, ok, err := c.UintArg(fieldName)
if err != nil {
return false, fmt.Errorf("reading SetRow() row: %v", err)
return false, fmt.Errorf("reading Store() row: %v", err)
} else if !ok {
return false, fmt.Errorf("SetRow() row argument '%v' required", rowLabel)
return false, fmt.Errorf("need the <FIELD>=<ROW> argument on Store()")
}
field := e.Holder.Field(index, fieldName)
@ -1774,7 +1793,7 @@ func (e *executor) executeSetRowShard(ctx context.Context, index string, c *pql.
}
src = row
} else {
return false, errors.New("SetRow() requires a source row")
return false, errors.New("Store() requires a source row")
}
// Set the row on the standard view.
@ -1793,7 +1812,7 @@ func (e *executor) executeSetRowShard(ctx context.Context, index string, c *pql.
}
set, err := fragment.setRow(src, rowID)
if err != nil {
return false, errors.Wrapf(err, "setting row %d on view %s shard %d", rowID, viewStandard, shard)
return false, errors.Wrapf(err, "storing row %d on view %s shard %d", rowID, viewStandard, shard)
}
changed = changed || set
@ -2354,7 +2373,7 @@ func (e *executor) translateCall(index string, idx *Index, c *pql.Call) error {
rowKey = "_" + rowLabel
fieldName = callArgString(c, "_field")
case "Rows":
fieldName = callArgString(c, "field")
fieldName = callArgString(c, "_field")
rowKey = "previous"
colKey = "column"
case "GroupBy":
@ -2459,7 +2478,7 @@ func (e *executor) translateGroupByCall(index string, idx *Index, c *pql.Call) e
fields := make([]*Field, len(c.Children))
for i, child := range c.Children {
fieldname := callArgString(child, "field")
fieldname := callArgString(child, "_field")
field := idx.Field(fieldname)
if field == nil {
return errors.Wrapf(ErrFieldNotFound, "getting field '%s' from '%s'", fieldname, child)
@ -2570,7 +2589,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res
case RowIDs:
other := RowIdentifiers{}
fieldName := callArgString(call, "field")
fieldName := callArgString(call, "_field")
if fieldName == "" {
return nil, ErrFieldNotFound
}
@ -2768,11 +2787,12 @@ func newGroupByIterator(rowIDs []RowIDs, children []*pql.Call, filter *Row, inde
fields: make([]FieldRow, len(children)),
}
var fieldName string
var ok bool
ignorePrev := false
for i, call := range children {
fieldName, ok := call.Args["field"].(string)
if !ok {
return nil, errors.Errorf("%s call must have 'field' argument with valid (string) field name. Got %v of type %[2]T", call.Name, call.Args["field"])
if fieldName, ok = call.Args["_field"].(string); !ok {
return nil, errors.Errorf("%s call must have field with valid (string) field name. Got %v of type %[2]T", call.Name, call.Args["_field"])
}
if holder.Field(index, fieldName) == nil {
return nil, ErrFieldNotFound

View file

@ -14,7 +14,7 @@ func TestExecutor_TranslateGroupByCall(t *testing.T) {
e := &executor{
Holder: NewHolder(),
}
e.Holder.Path, _ = ioutil.TempDir("", "")
e.Holder.Path, _ = ioutil.TempDir(*TempDir, "")
err := e.Holder.Open()
if err != nil {
t.Fatalf("opening holder: %v", err)
@ -46,7 +46,7 @@ func TestExecutor_TranslateGroupByCall(t *testing.T) {
t.Fatalf("translating rows %v, %v", erra, errb)
}
query, err := pql.ParseString(`GroupBy(Rows(field=ak), Rows(field=b), Rows(field=ck), previous=["la", 0, "ha"])`)
query, err := pql.ParseString(`GroupBy(Rows(ak), Rows(b), Rows(ck), previous=["la", 0, "ha"])`)
if err != nil {
t.Fatalf("parsing query: %v", err)
}
@ -69,28 +69,28 @@ func TestExecutor_TranslateGroupByCall(t *testing.T) {
err string
}{
{
pql: `GroupBy(Rows(field=notfound), previous=1)`,
pql: `GroupBy(Rows(notfound), previous=1)`,
err: "'previous' argument must be list",
},
{
pql: `GroupBy(Rows(field=ak), previous=["la", 0])`,
pql: `GroupBy(Rows(ak), previous=["la", 0])`,
err: "mismatched lengths",
},
{
pql: `GroupBy(Rows(field=ak), previous=[1])`,
pql: `GroupBy(Rows(ak), previous=[1])`,
err: "prev value must be a string",
},
{
pql: `GroupBy(Rows(field=notfound), previous=[1])`,
pql: `GroupBy(Rows(notfound), previous=[1])`,
err: ErrFieldNotFound.Error(),
},
// TODO: an unknown key will actually allocate an id. this is probably bad.
// {
// pql: `GroupBy(Rows(field=ak), previous=["zoop"])`,
// pql: `GroupBy(Rows(ak), previous=["zoop"])`,
// err: "translating row key '",
// },
{
pql: `GroupBy(Rows(field=b), previous=["la"])`,
pql: `GroupBy(Rows(b), previous=["la"])`,
err: "which doesn't use string keys",
},
}

View file

@ -16,7 +16,9 @@ package pilosa_test
import (
"context"
"flag"
"fmt"
"io/ioutil"
"math/rand"
"reflect"
"strconv"
@ -33,6 +35,22 @@ import (
"github.com/pkg/errors"
)
var (
TempDir = getTempDirString()
)
func getTempDirString() (td *string) {
tdflag := flag.Lookup("temp-dir")
if tdflag == nil {
td = flag.String("temp-dir", "", "Directory in which to place temporary data (e.g. for benchmarking). Useful if you are trying to benchmark different storage configurations.")
} else {
s := tdflag.Value.String()
td = &s
}
return td
}
// Ensure a row query can be executed.
func TestExecutor_Execute_Row(t *testing.T) {
t.Run("RowIDColumnID", func(t *testing.T) {
@ -2283,7 +2301,7 @@ Set(4500001, fn=4)
t.Run("remote groupBy", func(t *testing.T) {
if res, err := c[1].API.Query(context.Background(), &pilosa.QueryRequest{
Index: "i",
Query: `GroupBy(Rows(field=f))`,
Query: `GroupBy(Rows(f))`,
}); err != nil {
t.Fatalf("GroupBy querying: %v", err)
} else {
@ -2987,7 +3005,16 @@ func TestExecutor_Execute_SetRow(t *testing.T) {
}
func benchmarkExistence(nn bool, b *testing.B) {
c := test.MustRunCluster(b, 1)
c := test.MustNewCluster(b, 1)
var err error
c[0].Config.DataDir, err = ioutil.TempDir(*TempDir, "benchmarkExistence")
if err != nil {
b.Fatalf("getting temp dir: %v", err)
}
err = c.Start()
if err != nil {
b.Fatalf("starting cluster: %v", err)
}
defer c.Close()
hldr := test.Holder{Holder: c[0].Server.Holder()}
@ -3038,27 +3065,45 @@ func TestExecutor_Execute_Rows(t *testing.T) {
{13, 3},
})
rows := c.Query(t, "i", `Rows(field=general)`).Results[0].(pilosa.RowIdentifiers)
rows := c.Query(t, "i", `Rows(general)`).Results[0].(pilosa.RowIdentifiers)
if !reflect.DeepEqual(rows, pilosa.RowIdentifiers{Rows: []uint64{10, 11, 12, 13}}) {
t.Fatalf("unexpected rows: %+v", rows)
}
rows = c.Query(t, "i", `Rows(field=general, limit=2)`).Results[0].(pilosa.RowIdentifiers)
// backwards compatibility
// TODO: remove at Pilosa 2.0
rows = c.Query(t, "i", `Rows(field=general)`).Results[0].(pilosa.RowIdentifiers)
if !reflect.DeepEqual(rows, pilosa.RowIdentifiers{Rows: []uint64{10, 11, 12, 13}}) {
t.Fatalf("unexpected rows: %+v", rows)
}
rows = c.Query(t, "i", `Rows(general, limit=2)`).Results[0].(pilosa.RowIdentifiers)
if !reflect.DeepEqual(rows, pilosa.RowIdentifiers{Rows: []uint64{10, 11}}) {
t.Fatalf("unexpected rows: %+v", rows)
}
rows = c.Query(t, "i", `Rows(field=general, previous=10,limit=2)`).Results[0].(pilosa.RowIdentifiers)
rows = c.Query(t, "i", `Rows(general, previous=10,limit=2)`).Results[0].(pilosa.RowIdentifiers)
if !reflect.DeepEqual(rows, pilosa.RowIdentifiers{Rows: []uint64{11, 12}}) {
t.Fatalf("unexpected rows: %+v", rows)
}
rows = c.Query(t, "i", `Rows(field=general, column=2)`).Results[0].(pilosa.RowIdentifiers)
rows = c.Query(t, "i", `Rows(general, column=2)`).Results[0].(pilosa.RowIdentifiers)
if !reflect.DeepEqual(rows, pilosa.RowIdentifiers{Rows: []uint64{11, 12}}) {
t.Fatalf("unexpected rows: %+v", rows)
}
}
func TestExecutor_Execute_RowsTime(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
c.CreateField(t, "i", pilosa.IndexOptions{}, "t", pilosa.OptFieldTypeTime(pilosa.TimeQuantum("YMD"), true))
exp := "executing: Rows() query on time field with no standard view is not currently supported"
if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Rows(field=t)`}); err == nil || err.Error() != exp {
t.Fatalf("expected error: %s", exp)
}
}
func TestExecutor_Execute_Query_Error(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
@ -3070,35 +3115,27 @@ func TestExecutor_Execute_Query_Error(t *testing.T) {
}{
{
query: "GroupBy(Rows())",
error: "Rows call must have 'field' argument",
error: "Rows call must have field",
},
{
query: "GroupBy(Rows(field=true))",
error: "Rows call must have 'field' argument",
query: "GroupBy(Rows(\"true\"))",
error: "parsing: parsing:",
},
{
query: "GroupBy(Rows(field=\"true\"))",
error: "field not found",
query: "GroupBy(Rows(1))",
error: "parsing: parsing:",
},
{
query: "GroupBy(Rows(field=1))",
error: "Rows call must have 'field' argument",
},
{
query: "GroupBy(Rows(field))",
error: "parse error",
},
{
query: "GroupBy(Rows(field=general, limit=-1))",
query: "GroupBy(Rows(general, limit=-1))",
error: "must be positive, but got",
},
{
query: "GroupBy(Rows(field=general), limit=-1)",
query: "GroupBy(Rows(general), limit=-1)",
error: "must be positive, but got",
},
{
query: "GroupBy(Rows(field=general), filter=Rows(field=general))",
error: "unknown call: Rows",
query: "GroupBy(Rows(general), filter=Rows(general))",
error: "parsing: parsing:",
},
}
@ -3112,7 +3149,7 @@ func TestExecutor_Execute_Query_Error(t *testing.T) {
t.Fatalf("should have gotten an error on invalid rows query, but got %#v", r)
}
if !strings.Contains(err.Error(), test.error) {
t.Fatalf("unexpected error message: %s", err.Error())
t.Fatalf("unexpected error message:\n%s != %s", test.error, err.Error())
}
})
}
@ -3158,68 +3195,80 @@ func TestExecutor_Execute_Rows_Keys(t *testing.T) {
q string
exp []string
}{
{
q: `Rows(f)`,
exp: []string{"0", "1", "2", "3", "4", "5", "6", "7", "8", "9", "10", "11", "12", "13", "14", "15", "16", "17", "18"},
},
// backwards compatibility
// TODO: remove at Pilosa 2.0
{
q: `Rows(field=f)`,
exp: []string{"0", "1", "2", "3", "4", "5", "6", "7", "8", "9", "10", "11", "12", "13", "14", "15", "16", "17", "18"},
},
{
q: `Rows(f, limit=2)`,
exp: []string{"0", "1"},
},
// backwards compatibility
// TODO: remove at Pilosa 2.0
{
q: `Rows(field=f, limit=2)`,
exp: []string{"0", "1"},
},
{
q: `Rows(field=f, previous="15")`,
q: `Rows(f, previous="15")`,
exp: []string{"16", "17", "18"},
},
{
q: `Rows(field=f, previous="11", limit=2)`,
q: `Rows(f, previous="11", limit=2)`,
exp: []string{"12", "13"},
},
{
q: `Rows(field=f, previous="17", limit=5)`,
q: `Rows(f, previous="17", limit=5)`,
exp: []string{"18"},
},
{
q: `Rows(field=f, previous="18")`,
q: `Rows(f, previous="18")`,
exp: []string{},
},
{
q: `Rows(field=f, previous="1", limit=0)`,
q: `Rows(f, previous="1", limit=0)`,
exp: []string{},
},
{
q: `Rows(field=f, column="1")`,
q: `Rows(f, column="1")`,
exp: []string{"0", "1"},
},
{
q: `Rows(field=f, column="2")`,
q: `Rows(f, column="2")`,
exp: []string{"0", "1", "2"},
},
{
q: `Rows(field=f, column="3")`,
q: `Rows(f, column="3")`,
exp: []string{"1", "2", "3"},
},
{
q: `Rows(field=f, limit=2, column="3")`,
q: `Rows(f, limit=2, column="3")`,
exp: []string{"1", "2"},
},
{
q: fmt.Sprintf(`Rows(field=f, previous="15", column="%d")`, ShardWidth*9+17),
q: fmt.Sprintf(`Rows(f, previous="15", column="%d")`, ShardWidth*9+17),
exp: []string{"16", "17"},
},
{
q: fmt.Sprintf(`Rows(field=f, previous="11", limit=2, column="%d")`, ShardWidth*5+14),
q: fmt.Sprintf(`Rows(f, previous="11", limit=2, column="%d")`, ShardWidth*5+14),
exp: []string{"12", "13"},
},
{
q: fmt.Sprintf(`Rows(field=f, previous="17", limit=5, column="%d")`, ShardWidth*9+18),
q: fmt.Sprintf(`Rows(f, previous="17", limit=5, column="%d")`, ShardWidth*9+18),
exp: []string{"18"},
},
{
q: `Rows(field=f, previous="18", column="19")`,
q: `Rows(f, previous="18", column="19")`,
exp: []string{},
},
{
q: `Rows(field=f, previous="1", limit=0, column="0")`,
q: `Rows(f, previous="1", limit=0, column="0")`,
exp: []string{},
},
}
@ -3271,13 +3320,27 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
})
t.Run("Unknown Field ", func(t *testing.T) {
if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `GroupBy(Rows(field=missing))`}); err != nil {
if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `GroupBy(Rows(missing))`}); err != nil {
if errors.Cause(err) != pilosa.ErrFieldNotFound {
t.Fatalf("unexpected error\n\"%s\" not returned instead \n\"%s\"", pilosa.ErrFieldNotFound, err)
}
}
})
// backwards compatibility
// TODO: remove at Pilosa 2.0
t.Run("BasicLegacy", func(t *testing.T) {
expected := []pilosa.GroupCount{
{Group: []pilosa.FieldRow{{Field: "general", RowID: 10}, {Field: "sub", RowID: 100}}, Count: 3},
{Group: []pilosa.FieldRow{{Field: "general", RowID: 10}, {Field: "sub", RowID: 110}}, Count: 1},
{Group: []pilosa.FieldRow{{Field: "general", RowID: 11}, {Field: "sub", RowID: 110}}, Count: 1},
{Group: []pilosa.FieldRow{{Field: "general", RowID: 12}, {Field: "sub", RowID: 110}}, Count: 1},
}
results := c.Query(t, "i", `GroupBy(Rows(field=general), Rows(sub))`).Results[0].([]pilosa.GroupCount)
test.CheckGroupBy(t, expected, results)
})
t.Run("Basic", func(t *testing.T) {
expected := []pilosa.GroupCount{
{Group: []pilosa.FieldRow{{Field: "general", RowID: 10}, {Field: "sub", RowID: 100}}, Count: 3},
@ -3286,7 +3349,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
{Group: []pilosa.FieldRow{{Field: "general", RowID: 12}, {Field: "sub", RowID: 110}}, Count: 1},
}
results := c.Query(t, "i", `GroupBy(Rows(field=general), Rows(field=sub))`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(general), Rows(sub))`).Results[0].([]pilosa.GroupCount)
test.CheckGroupBy(t, expected, results)
})
@ -3296,7 +3359,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
{Group: []pilosa.FieldRow{{Field: "general", RowID: 10}, {Field: "sub", RowID: 110}}, Count: 1},
}
results := c.Query(t, "i", `GroupBy(Rows(field=general), Rows(field=sub), filter=Row(general=10))`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(general), Rows(sub), filter=Row(general=10))`).Results[0].([]pilosa.GroupCount)
test.CheckGroupBy(t, expected, results)
})
@ -3306,7 +3369,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
{Group: []pilosa.FieldRow{{Field: "general", RowID: 12}}, Count: 2},
}
results := c.Query(t, "i", `GroupBy(Rows(field=general, previous=10))`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(general, previous=10))`).Results[0].([]pilosa.GroupCount)
test.CheckGroupBy(t, expected, results)
})
@ -3315,7 +3378,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
{Group: []pilosa.FieldRow{{Field: "general", RowID: 11}}, Count: 2},
}
results := c.Query(t, "i", `GroupBy(Rows(field=general, previous=10), limit=1)`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(general, previous=10), limit=1)`).Results[0].([]pilosa.GroupCount)
test.CheckGroupBy(t, expected, results)
})
@ -3336,7 +3399,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
{Group: []pilosa.FieldRow{{Field: "a", RowID: 0}, {Field: "b", RowID: 1}}, Count: 1},
}
results := c.Query(t, "i", `GroupBy(Rows(field=a), Rows(field=b), limit=1)`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(a), Rows(b), limit=1)`).Results[0].([]pilosa.GroupCount)
test.CheckGroupBy(t, expected, results)
})
@ -3364,7 +3427,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
})
t.Run("test wrapping with previous", func(t *testing.T) {
results := c.Query(t, "i", `GroupBy(Rows(field=wa), Rows(field=wb), Rows(field=wc, previous=1), limit=3)`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(wa), Rows(wb), Rows(wc, previous=1), limit=3)`).Results[0].([]pilosa.GroupCount)
expected := []pilosa.GroupCount{
{Group: []pilosa.FieldRow{{Field: "wa", RowID: 0}, {Field: "wb", RowID: 0}, {Field: "wc", RowID: 2}}, Count: 2},
{Group: []pilosa.FieldRow{{Field: "wa", RowID: 0}, {Field: "wb", RowID: 1}, {Field: "wc", RowID: 0}}, Count: 1},
@ -3374,14 +3437,14 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
})
t.Run("test previous is last result", func(t *testing.T) {
results := c.Query(t, "i", `GroupBy(Rows(field=wa, previous=3), Rows(field=wb, previous=3), Rows(field=wc, previous=3), limit=3)`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(wa, previous=3), Rows(wb, previous=3), Rows(wc, previous=3), limit=3)`).Results[0].([]pilosa.GroupCount)
if len(results) > 0 {
t.Fatalf("expected no results because previous specified last result")
}
})
t.Run("test wrapping multiple", func(t *testing.T) {
results := c.Query(t, "i", `GroupBy(Rows(field=wa), Rows(field=wb, previous=2), Rows(field=wc, previous=2), limit=1)`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(wa), Rows(wb, previous=2), Rows(wc, previous=2), limit=1)`).Results[0].([]pilosa.GroupCount)
expected := []pilosa.GroupCount{
{Group: []pilosa.FieldRow{{Field: "wa", RowID: 1}, {Field: "wb", RowID: 0}, {Field: "wc", RowID: 0}}, Count: 1},
}
@ -3405,7 +3468,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
{3, ShardWidth},
})
t.Run("distinct rows in different shards", func(t *testing.T) {
results := c.Query(t, "i", `GroupBy(Rows(field=ma), Rows(field=mb), limit=5)`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(ma), Rows(mb), limit=5)`).Results[0].([]pilosa.GroupCount)
expected := []pilosa.GroupCount{
{Group: []pilosa.FieldRow{{Field: "ma", RowID: 0}, {Field: "mb", RowID: 0}}, Count: 1},
{Group: []pilosa.FieldRow{{Field: "ma", RowID: 0}, {Field: "mb", RowID: 2}}, Count: 1},
@ -3417,7 +3480,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
})
t.Run("distinct rows in different shards with row limit", func(t *testing.T) {
results := c.Query(t, "i", `GroupBy(Rows(field=ma), Rows(field=mb, limit=2), limit=5)`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(ma), Rows(mb, limit=2), limit=5)`).Results[0].([]pilosa.GroupCount)
expected := []pilosa.GroupCount{
{Group: []pilosa.FieldRow{{Field: "ma", RowID: 0}, {Field: "mb", RowID: 0}}, Count: 1},
{Group: []pilosa.FieldRow{{Field: "ma", RowID: 1}, {Field: "mb", RowID: 1}}, Count: 1},
@ -3428,7 +3491,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
})
t.Run("distinct rows in different shards with column arg", func(t *testing.T) {
results := c.Query(t, "i", fmt.Sprintf(`GroupBy(Rows(field=ma), Rows(field=mb, column=%d), limit=5)`, ShardWidth)).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", fmt.Sprintf(`GroupBy(Rows(ma), Rows(mb, column=%d), limit=5)`, ShardWidth)).Results[0].([]pilosa.GroupCount)
expected := []pilosa.GroupCount{
{Group: []pilosa.FieldRow{{Field: "ma", RowID: 1}, {Field: "mb", RowID: 1}}, Count: 1},
{Group: []pilosa.FieldRow{{Field: "ma", RowID: 1}, {Field: "mb", RowID: 3}}, Count: 1},
@ -3453,7 +3516,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
{1, ShardWidth},
})
t.Run("same rows in different shards", func(t *testing.T) {
results := c.Query(t, "i", `GroupBy(Rows(field=na), Rows(field=nb))`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(na), Rows(nb))`).Results[0].([]pilosa.GroupCount)
expected := []pilosa.GroupCount{
{Group: []pilosa.FieldRow{{Field: "na", RowID: 0}, {Field: "nb", RowID: 0}}, Count: 2},
{Group: []pilosa.FieldRow{{Field: "na", RowID: 0}, {Field: "nb", RowID: 1}}, Count: 2},
@ -3490,11 +3553,11 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
t.Run("test wrapping with previous", func(t *testing.T) {
totalResults := make([]pilosa.GroupCount, 0)
results := c.Query(t, "i", `GroupBy(Rows(field=ppa), Rows(field=ppb), Rows(field=ppc), limit=3)`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(ppa), Rows(ppb), Rows(ppc), limit=3)`).Results[0].([]pilosa.GroupCount)
totalResults = append(totalResults, results...)
for len(totalResults) < 64 {
lastGroup := results[len(results)-1].Group
query := fmt.Sprintf("GroupBy(Rows(field=ppa, previous=%d), Rows(field=ppb, previous=%d), Rows(field=ppc, previous=%d), limit=3)", lastGroup[0].RowID, lastGroup[1].RowID, lastGroup[2].RowID)
query := fmt.Sprintf("GroupBy(Rows(ppa, previous=%d), Rows(ppb, previous=%d), Rows(ppc, previous=%d), limit=3)", lastGroup[0].RowID, lastGroup[1].RowID, lastGroup[2].RowID)
results = c.Query(t, "i", query).Results[0].([]pilosa.GroupCount)
totalResults = append(totalResults, results...)
}
@ -3537,7 +3600,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
{Group: []pilosa.FieldRow{{Field: "generalk", RowID: 3, RowKey: "twelve"}, {Field: "subk", RowID: 2, RowKey: "one-hundred-ten"}}, Count: 1},
}
results := c.Query(t, "i", `GroupBy(Rows(field="generalk"), Rows(field="subk"))`).Results[0].([]pilosa.GroupCount)
results := c.Query(t, "i", `GroupBy(Rows(generalk), Rows(subk))`).Results[0].([]pilosa.GroupCount)
test.CheckGroupBy(t, expected, results)
})
@ -3551,7 +3614,16 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
}
func BenchmarkGroupBy(b *testing.B) {
c := test.MustRunCluster(b, 1)
c := test.MustNewCluster(b, 1)
var err error
c[0].Config.DataDir, err = ioutil.TempDir(*TempDir, "benchmarkGroupBy")
if err != nil {
b.Fatalf("getting temp dir: %v", err)
}
err = c.Start()
if err != nil {
b.Fatalf("starting cluster: %v", err)
}
defer c.Close()
c.CreateField(b, "i", pilosa.IndexOptions{}, "a")
c.CreateField(b, "i", pilosa.IndexOptions{}, "b")
@ -3586,7 +3658,7 @@ func BenchmarkGroupBy(b *testing.B) {
b.ResetTimer()
b.ReportAllocs()
for i := 0; i < b.N; i++ {
c.Query(b, "i", `GroupBy(Rows(field=a), Rows(field=b), Rows(field=c))`)
c.Query(b, "i", `GroupBy(Rows(a), Rows(b), Rows(c))`)
}
})
@ -3594,7 +3666,7 @@ func BenchmarkGroupBy(b *testing.B) {
b.ResetTimer()
b.ReportAllocs()
for i := 0; i < b.N; i++ {
c.Query(b, "i", `GroupBy(Rows(field=a), Rows(field=b), Rows(field=c), limit=4)`)
c.Query(b, "i", `GroupBy(Rows(a), Rows(b), Rows(c), limit=4)`)
}
})

View file

@ -192,7 +192,7 @@ type TestField struct {
// NewTestField returns a new instance of TestField d/0.
func NewTestField(opts FieldOption) *TestField {
path, err := ioutil.TempDir("", "pilosa-field-")
path, err := ioutil.TempDir(*TempDir, "pilosa-field-")
if err != nil {
panic(err)
}

View file

@ -1050,7 +1050,7 @@ func (f *fragment) top(opt topOptions) ([]Pair, error) {
rowID, cnt := pair.ID, pair.Count
// Ignore empty rows.
if cnt <= 0 {
if cnt == 0 {
continue
}

View file

@ -39,12 +39,9 @@ var (
// In order to generate the sample fragment file,
// run an import and copy PILOSA_DATA_DIR/INDEX_NAME/FRAME_NAME/0 to testdata/sample_view
FragmentPath = flag.String("fragment", "testdata/sample_view/0", "fragment path")
TempDir = ""
)
func init() { // nolint: gochecknoinits
flag.StringVar(&TempDir, "temp-dir", "", "Directory in which to place temporary data (e.g. for benchmarking). Useful if you are trying to benchmark different storage configurations.")
}
TempDir = flag.String("temp-dir", "", "Directory in which to place temporary data (e.g. for benchmarking). Useful if you are trying to benchmark different storage configurations.")
)
// Ensure a fragment can set a bit and retrieve it.
func TestFragment_SetBit(t *testing.T) {
@ -367,6 +364,20 @@ func TestFragment_Sum(t *testing.T) {
t.Fatalf("unexpected sum: %d", sum)
}
})
// verify that clearValue clears values
if _, err := f.clearValue(1000, bitDepth, 23); err != nil {
t.Fatal(err)
}
t.Run("ClearValue", func(t *testing.T) {
if sum, n, err := f.sum(nil, bitDepth); err != nil {
t.Fatal(err)
} else if n != 3 {
t.Fatalf("unexpected count: %d", n)
} else if sum != (3800 - 382) {
t.Fatalf("unexpected sum: got %d, expecting %d", sum, 3800-382)
}
})
}
// Ensure a fragment can find the min and max of values.
@ -637,6 +648,71 @@ func TestFragment_Range(t *testing.T) {
})
}
// benchmarkSetValues is a helper function to explore, very roughly, the cost
// of setting values.
func benchmarkSetValues(b *testing.B, bitDepth uint, f *fragment, cfunc func(uint64) uint64) {
column := uint64(0)
for i := 0; i < b.N; i++ {
f.setValue(column, bitDepth, uint64(i))
column = cfunc(column)
}
}
// Benchmark performance of setValue for BSI ranges.
func BenchmarkFragment_SetValue(b *testing.B) {
depths := []uint{4, 8, 16}
for _, bitDepth := range depths {
name := fmt.Sprintf("Depth%d", bitDepth)
f := mustOpenFragment("i", "f", viewBSIGroupPrefix+"foo", 0, "none")
b.Run(name+"_Sparse", func(b *testing.B) {
benchmarkSetValues(b, bitDepth, f, func(u uint64) uint64 { return (u + 70000) & (ShardWidth - 1) })
})
f.Clean(b)
f = mustOpenFragment("i", "f", viewBSIGroupPrefix+"foo", 0, "none")
b.Run(name+"_Dense", func(b *testing.B) {
benchmarkSetValues(b, bitDepth, f, func(u uint64) uint64 { return (u + 1) & (ShardWidth - 1) })
})
f.Clean(b)
}
}
// benchmarkImportValues is a helper function to explore, very roughly, the cost
// of setting values using the special setter used for imports.
func benchmarkImportValues(b *testing.B, bitDepth uint, f *fragment, cfunc func(uint64) uint64) {
column := uint64(0)
b.StopTimer()
columns := make([]uint64, b.N)
values := make([]uint64, b.N)
for i := 0; i < b.N; i++ {
values[i] = uint64(i)
columns[i] = column
column = cfunc(column)
}
b.StartTimer()
err := f.importValue(columns, values, bitDepth, false)
if err != nil {
b.Fatalf("error importing values: %s", err)
}
}
// Benchmark performance of setValue for BSI ranges.
func BenchmarkFragment_ImportValue(b *testing.B) {
depths := []uint{4, 8, 16}
for _, bitDepth := range depths {
name := fmt.Sprintf("Depth%d", bitDepth)
f := mustOpenFragment("i", "f", viewBSIGroupPrefix+"foo", 0, "none")
b.Run(name+"_Sparse", func(b *testing.B) {
benchmarkImportValues(b, bitDepth, f, func(u uint64) uint64 { return (u + 70000) & (ShardWidth - 1) })
})
f.Clean(b)
f = mustOpenFragment("i", "f", viewBSIGroupPrefix+"foo", 0, "none")
b.Run(name+"_Dense", func(b *testing.B) {
benchmarkImportValues(b, bitDepth, f, func(u uint64) uint64 { return (u + 1) & (ShardWidth - 1) })
})
f.Clean(b)
}
}
// Ensure a fragment can snapshot correctly.
func TestFragment_Snapshot(t *testing.T) {
f := mustOpenFragment("i", "f", viewStandard, 0, "")
@ -1142,7 +1218,7 @@ func BenchmarkFragment_Blocks(b *testing.B) {
if err := f.Open(); err != nil {
b.Fatal(err)
}
defer f.Clean(b)
defer f.CleanKeep(b)
// Reset timer and execute benchmark.
b.ResetTimer()
@ -1671,7 +1747,7 @@ func BenchmarkFragment_Snapshot(b *testing.B) {
if err := f.Open(); err != nil {
b.Fatal(err)
}
defer f.Clean(b)
defer f.CleanKeep(b)
b.ResetTimer()
// Reset timer and execute benchmark.
@ -1990,7 +2066,7 @@ func BenchmarkFileWrite(b *testing.B) {
b.Run(fmt.Sprintf("Rows%d", numRows), func(b *testing.B) {
b.StopTimer()
for i := 0; i < b.N; i++ {
f, err := ioutil.TempFile(TempDir, "")
f, err := ioutil.TempFile(*TempDir, "")
if err != nil {
b.Fatalf("getting temp file: %v", err)
}
@ -2030,9 +2106,23 @@ func (f *fragment) Clean(t testing.TB) {
}
}
// CleanKeep is just like Clean(), but it doesn't remove the
// fragment file (note that it DOES remove the cache file).
func (f *fragment) CleanKeep(t testing.TB) {
errc := f.Close()
errp := os.Remove(f.cachePath())
if errc != nil {
t.Fatal("closing fragment: ", errc, errp)
}
// not all fragments have cache files
if errp != nil && !os.IsNotExist(errp) {
t.Fatalf("cleaning up fragment cache: %v", errp)
}
}
// mustOpenFragment returns a new instance of Fragment with a temporary path.
func mustOpenFragment(index, field, view string, shard uint64, cacheType string) *fragment {
file, err := ioutil.TempFile(TempDir, "pilosa-fragment-")
file, err := ioutil.TempFile(*TempDir, "pilosa-fragment-")
if err != nil {
panic(err)
}

View file

@ -22,6 +22,7 @@ import (
"io/ioutil"
"log"
"net"
"os"
"strconv"
"strings"
"sync"
@ -50,8 +51,11 @@ type memberSet struct {
Logger logger.Logger
logger *log.Logger
// stdLogger is only used when passed into memberlist library things that take a std library logger rather than an interface.
stdLogger *log.Logger
// logOutput is similar to stdLogger in that it's passed to memberlist things which can't take a pilosa Logger.
logOutput io.Writer
transport *Transport
eventReceiver *eventReceiver
@ -151,14 +155,20 @@ func WithTransport(transport *Transport) memberSetOption {
}
}
// WithLogger is a functional option for providing a logger to NewMemberSet.
// WithLogger is a functional option for providing a Go logger to NewMemberSet.
// If the memberSet's transport is nil, this logger will be used when creating
// one. If WithLogOutput is not used, this logger will be passed to memberlist
// for it to use internally. This logger is not used for logging by code in this
// (gossip) package - for that, use the WithPilosaLogger option.
func WithLogger(logger *log.Logger) memberSetOption {
return func(g *memberSet) error {
g.logger = logger
g.stdLogger = logger
return nil
}
}
// WithLogOutput allows one to pass a Writer which will in turn be passed to
// memberlist for use in logging.
func WithLogOutput(o io.Writer) memberSetOption {
return func(g *memberSet) error {
g.logOutput = o
@ -166,7 +176,20 @@ func WithLogOutput(o io.Writer) memberSetOption {
}
}
// NewMemberSet returns a new instance of GossipMemberSet based on options.
// WithPilosaLogger allows one to configure a memberSet with a logger of their
// choice which satisfies the pilosa logger interface.
func WithPilosaLogger(l logger.Logger) memberSetOption {
return func(g *memberSet) error {
g.Logger = l
return nil
}
}
// NewMemberSet returns a new instance of GossipMemberSet based on options. The
// logging options which can be passed to NewMemberSet are complicated for
// historical reasons - please pass WithPilosaLogger, and either WithLogOutput
// or WithLogger. If you pass WithLogOutput, be sure to also pass in a Transport
// using WithTransport.
func NewMemberSet(cfg Config, api *pilosa.API, options ...memberSetOption) (*memberSet, error) {
host := api.Node().URI.Host
g := &memberSet{
@ -180,7 +203,8 @@ func NewMemberSet(cfg Config, api *pilosa.API, options ...memberSetOption) (*mem
return nil, errors.Wrap(err, "executing option")
}
}
ger := newEventReceiver(g.logger, api)
ger := newEventReceiver(g.Logger, api)
g.eventReceiver = ger
if g.transport == nil {
@ -189,8 +213,16 @@ func NewMemberSet(cfg Config, api *pilosa.API, options ...memberSetOption) (*mem
return nil, fmt.Errorf("convert port: %s", err)
}
if g.stdLogger == nil {
if g.logOutput != nil {
g.stdLogger = logger.NewStandardLogger(g.logOutput).Logger()
} else {
g.stdLogger = log.New(os.Stderr, "", log.LstdFlags)
}
}
// Set up the transport.
transport, err := NewTransport(host, port, g.logger)
transport, err := NewTransport(host, port, g.stdLogger)
if err != nil {
return nil, fmt.Errorf("new tranport: %s", err)
}
@ -233,7 +265,7 @@ func NewMemberSet(cfg Config, api *pilosa.API, options ...memberSetOption) (*mem
if g.logOutput != nil {
conf.LogOutput = g.logOutput
} else {
conf.Logger = g.logger
conf.Logger = g.stdLogger
}
g.config = &config{
@ -318,11 +350,11 @@ type eventReceiver struct {
ch chan memberlist.NodeEvent
papi *pilosa.API
logger *log.Logger
logger logger.Logger
}
// newEventReceiver returns a new instance of GossipEventReceiver.
func newEventReceiver(logger *log.Logger, papi *pilosa.API) *eventReceiver {
func newEventReceiver(logger logger.Logger, papi *pilosa.API) *eventReceiver {
ger := &eventReceiver{
ch: make(chan memberlist.NodeEvent, 1),
logger: logger,
@ -468,7 +500,7 @@ func newTransport(conf *memberlist.Config) (*memberlist.NetTransport, error) {
nt, err := makeNetRetry(limit)
if err != nil {
return nil, fmt.Errorf("Could not set up network transport: %v", err)
return nil, errors.Wrap(err, "could not set up network transport")
}
return nt, nil

View file

@ -46,7 +46,7 @@ func (h *tHolder) Reopen() error {
}
func newHolder() *tHolder {
path, err := ioutil.TempDir("", "pilosa-")
path, err := ioutil.TempDir(*TempDir, "pilosa-")
if err != nil {
panic(err)
}

View file

@ -584,11 +584,11 @@ func validateOptions(data map[string]interface{}, validIndexOptions []string) er
}
for kk, vv := range options {
if !foundItem(validIndexOptions, kk) {
return fmt.Errorf("Unknown key: %v:%v", kk, vv)
return fmt.Errorf("unknown key: %v:%v", kk, vv)
}
}
default:
return fmt.Errorf("Unknown key: %v:%v", k, v)
return fmt.Errorf("unknown key: %v:%v", k, v)
}
}
return nil

View file

@ -35,8 +35,8 @@ func TestPostIndexRequestUnmarshalJSON(t *testing.T) {
{json: `{"options": {"trackExistence": false}}`, expected: postIndexRequest{Options: pilosa.IndexOptions{TrackExistence: false}}},
{json: `{"options": {"keys": true}}`, expected: postIndexRequest{Options: pilosa.IndexOptions{Keys: true, TrackExistence: true}}},
{json: `{"options": 4}`, err: "options is not map[string]interface{}"},
{json: `{"option": {}}`, err: "Unknown key: option:map[]"},
{json: `{"options": {"badKey": "test"}}`, err: "Unknown key: badKey:test"},
{json: `{"option": {}}`, err: "unknown key: option:map[]"},
{json: `{"options": {"badKey": "test"}}`, err: "unknown key: badKey:test"},
}
for _, test := range tests {
actual := &postIndexRequest{}

View file

@ -21,7 +21,7 @@ import (
// mustOpenIndex returns a new, opened index at a temporary path. Panic on error.
func mustOpenIndex(opt IndexOptions) *Index {
path, err := ioutil.TempDir("", "pilosa-index-")
path, err := ioutil.TempDir(*TempDir, "pilosa-index-")
if err != nil {
panic(err)
}

View file

@ -83,7 +83,7 @@ func (c *Cache) Get(key Key) (value interface{}, ok bool) {
}
// remove removes the provided key from the cache.
func (c *Cache) remove(key Key) { // nolint: megacheck
func (c *Cache) remove(key Key) { // nolint: staticcheck
if c.cache == nil {
return
}
@ -121,7 +121,7 @@ func (c *Cache) Len() int {
}
// clear purges all stored items from the cache.
func (c *Cache) clear() { // nolint: megacheck
func (c *Cache) clear() { // nolint: staticcheck
if c.OnEvicted != nil {
for _, e := range c.cache {
kv := e.Value.(*entry)

View file

@ -259,7 +259,7 @@ func (c *Call) FieldArg() (string, error) {
return arg, nil
}
}
return "", fmt.Errorf("No field argument specified")
return "", fmt.Errorf("no field argument specified")
}
func IsReservedArg(name string) bool {

View file

@ -13,6 +13,7 @@ Call <- 'Set' {p.startCall("Set")} open col comma args (comma timestamp)? close
/ 'ClearRow' {p.startCall("ClearRow")} open arg close {p.endCall()}
/ 'Store' {p.startCall("Store")} open Call comma arg close {p.endCall()}
/ 'TopN' {p.startCall("TopN")} open posfield (comma allargs)? close {p.endCall()}
/ 'Rows' {p.startCall("Rows")} open posfield (comma allargs)? close {p.endCall()}
/ < IDENT > { p.startCall(buffer[begin:end] ) } open allargs comma? close { p.endCall() }
allargs <- Call (comma Call)* (comma args)? / args / sp
args <- arg (comma args)? sp

File diff suppressed because it is too large Load diff

View file

@ -163,9 +163,8 @@ func (b *Bitmap) Add(a ...uint64) (changed bool, err error) {
}
// Apply to the in-memory bitmap.
if op.apply(b) {
if b.DirectAdd(v) {
changed = true
}
}
@ -233,6 +232,18 @@ func (b *Bitmap) Count() (n uint64) {
return b.Containers.Count()
}
// Size returns the number of bytes required for the bitmap.
func (b *Bitmap) Size() int {
numbytes := 0
citer, _ := b.Containers.Iterator(0)
for citer.Next() {
_, c := citer.Value()
numbytes += c.size()
}
return numbytes
}
// CountRange returns the number of bits set between [start, end).
func (b *Bitmap) CountRange(start, end uint64) (n uint64) {
if b.Containers.Size() == 0 {
@ -3973,7 +3984,7 @@ func readOfficialHeader(buf []byte) (size uint32, containerTyper func(index uint
header = pos
if size > (1 << 16) {
err = fmt.Errorf("It is logically impossible to have more than (1<<16) containers.")
err = fmt.Errorf("it is logically impossible to have more than (1<<16) containers")
return size, containerTyper, header, pos, haveRuns, err
}

View file

@ -3224,11 +3224,11 @@ func TestContainerCombinations(t *testing.T) {
//func getFunc(func(a, b *container) *container, m, n *container) *container {
func runContainerFunc(f interface{}, c ...*Container) *Container {
switch f.(type) {
switch f := f.(type) {
case func(*Container) *Container:
return f.(func(*Container) *Container)(c[0])
return f(c[0])
case func(*Container, *Container) *Container:
return f.(func(a, b *Container) *Container)(c[0], c[1])
return f(c[0], c[1])
}
return nil
}

View file

@ -37,6 +37,29 @@ func TestContainerCount(t *testing.T) {
t.Fatalf("Count != CountRange\n")
}
}
func TestSize(t *testing.T) {
//array
a := roaring.NewFileBitmap(0, 65535, 131072)
if a.Size() != 6 {
t.Fatalf("Size in bytes incorrect \n")
}
//bitmap
b := roaring.NewFileBitmap()
for i := uint64(0); i <= 4096; i++ {
b.DirectAdd(i)
}
if b.Size() != 8192 {
t.Fatalf("Size in bytes incorrect \n")
}
//convert to rle
b.Optimize()
//rle
if b.Size() != 6 {
t.Fatalf("Size in bytes incorrect \n")
}
}
func TestCountRange(t *testing.T) {
tests := []struct {
@ -1498,6 +1521,41 @@ func BenchmarkSliceDescending(b *testing.B) {
for col := uint64(pilosa.ShardWidth); col > uint64(0); col-- {
bm.Add(col)
}
bm.Add(0)
}
}
func BenchmarkSliceAscendingStriped(b *testing.B) {
for n := 0; n < b.N; n++ {
bm := roaring.NewFileBitmap()
l := uint64(pilosa.ShardWidth / 8)
for col := uint64(0); col < l; col++ {
bm.Add(l*0 + col)
bm.Add(l*1 + col)
bm.Add(l*2 + col)
bm.Add(l*3 + col)
bm.Add(l*4 + col)
bm.Add(l*5 + col)
bm.Add(l*6 + col)
bm.Add(l*7 + col)
}
}
}
func BenchmarkSliceDescendingStriped(b *testing.B) {
for n := 0; n < b.N; n++ {
bm := roaring.NewFileBitmap()
l := uint64(pilosa.ShardWidth / 8)
for col := uint64(l); col < l+1; col-- {
bm.Add(l*7 + col)
bm.Add(l*6 + col)
bm.Add(l*5 + col)
bm.Add(l*4 + col)
bm.Add(l*3 + col)
bm.Add(l*2 + col)
bm.Add(l*1 + col)
bm.Add(l*0 + col)
}
}
}

View file

@ -487,7 +487,7 @@ func (s *Server) receiveMessage(m Message) error {
case *CreateShardMessage:
f := s.holder.Field(obj.Index, obj.Field)
if f == nil {
return fmt.Errorf("Local field not found: %s/%s", obj.Index, obj.Field)
return fmt.Errorf("local field not found: %s/%s", obj.Index, obj.Field)
}
if err := f.AddRemoteAvailableShards(roaring.NewBitmap(obj.Shard)); err != nil {
return errors.Wrap(err, "adding remote available shards")
@ -505,7 +505,7 @@ func (s *Server) receiveMessage(m Message) error {
case *CreateFieldMessage:
idx := s.holder.Index(obj.Index)
if idx == nil {
return fmt.Errorf("Local Index not found: %s", obj.Index)
return fmt.Errorf("local index not found: %s", obj.Index)
}
opt := obj.Meta
_, err := idx.createField(obj.Field, *opt)
@ -525,7 +525,7 @@ func (s *Server) receiveMessage(m Message) error {
case *CreateViewMessage:
f := s.holder.Field(obj.Index, obj.Field)
if f == nil {
return fmt.Errorf("Local Field not found: %s", obj.Field)
return fmt.Errorf("local field not found: %s", obj.Field)
}
_, _, err := f.createViewIfNotExistsBase(obj.View)
if err != nil {
@ -534,7 +534,7 @@ func (s *Server) receiveMessage(m Message) error {
case *DeleteViewMessage:
f := s.holder.Field(obj.Index, obj.Field)
if f == nil {
return fmt.Errorf("Local Field not found: %s", obj.Field)
return fmt.Errorf("local field not found: %s", obj.Field)
}
err := f.deleteView(obj.View)
if err != nil {

View file

@ -399,7 +399,7 @@ func TestClusterResize_RemoveNode(t *testing.T) {
nodeID := mustNodeID(m0.URL())
resp := test.MustDo("POST", m0.URL()+fmt.Sprintf("/cluster/resize/remove-node"), fmt.Sprintf(`{"id": "%s"}`, nodeID))
expBody := "removing node: calling node leave: coordinator 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 {

View file

@ -320,6 +320,7 @@ func (m *Command) setupNetworking() error {
m.Config.Gossip,
m.API,
gossip.WithLogOutput(&filteredWriter{logOutput: m.logOutput, v: m.Config.Verbose}),
gossip.WithPilosaLogger(m.logger),
gossip.WithTransport(m.gossipTransport),
)
if err != nil {

View file

@ -293,7 +293,7 @@ func TestMain_GroupBy(t *testing.T) {
}
// Query row.
if res, err := m.QueryProtobuf("i", `GroupBy(Rows(field="generalk"), Rows(field="subk"))`); err != nil {
if res, err := m.QueryProtobuf("i", `GroupBy(Rows(generalk), Rows(subk))`); err != nil {
t.Fatal(err)
} else {
test.CheckGroupBy(t, expected, res.Results[0].([]pilosa.GroupCount))

View file

@ -38,7 +38,7 @@ func TestCountOpenFiles(t *testing.T) {
func TestMonitorAntiEntropyZero(t *testing.T) {
td, err := ioutil.TempDir("", "")
td, err := ioutil.TempDir(*TempDir, "")
if err != nil {
t.Fatalf("getting temp dir: %v", err)
}

View file

@ -372,7 +372,6 @@ func (s *TranslateFile) monitorReplication() {
if err := s.replicate(ctx); err != nil {
s.logger.Printf("pilosa: replication error: %s", err)
}
select {
case <-ctx.Done():
return
@ -412,22 +411,42 @@ func (s *TranslateFile) replicate(ctx context.Context) error {
// Wrap in bufferred I/O so it implements io.ByteReader.
bufr := bufio.NewReader(rc)
// we need a way to make an asynchronous routine hand us back an error,
// but we might not still be there to get it. so we have a buffer.
chErr := make(chan error, 1)
// Continually read new entries from primary and append to local store.
for {
// Read next available entry.
var entry LogEntry
if _, err := entry.ReadFrom(bufr); err == io.EOF {
if _, err = entry.ReadFrom(bufr); err == io.EOF {
return nil
} else if err != nil {
return err
}
s.mu.Lock()
// Write to local store.
if err := s.appendEntry(&entry); err != nil {
s.mu.Unlock()
return err
// note: we should never end up spawning two of this goroutine
// at once. either we end up reading the error from chErr below,
// and this loop continues, or we don't, and the whole function
// returns. if the function returns, we can write that single
// error to the empty channel with a buffer of 1, the goroutine
// terminates, and chErr becomes garbage-collectable.
go func() {
s.mu.Lock()
defer s.mu.Unlock()
// Write to local store.
err = s.appendEntry(&entry)
chErr <- err
}()
select {
case err = <-chErr:
if err != nil {
return err
}
case <-s.replicationClosing:
return nil
case <-ctx.Done():
return nil
}
s.mu.Unlock()
}
}

View file

@ -803,7 +803,7 @@ type TranslateFile struct {
}
func NewTranslateFile() *TranslateFile {
f, err := ioutil.TempFile("", "")
f, err := ioutil.TempFile(*TempDir, "")
if err != nil {
panic(err)
}

View file

@ -212,7 +212,7 @@ func (t *ClusterCluster) addCluster(i int, saveTopology bool) (*cluster, error)
t.common.Nodes = append(t.common.Nodes, node)
// create node-specific temp directory
path, err := ioutil.TempDir("", fmt.Sprintf("pilosa-cluster-node-%d-", i))
path, err := ioutil.TempDir(*TempDir, fmt.Sprintf("pilosa-cluster-node-%d-", i))
if err != nil {
return nil, err
}

View file

@ -24,7 +24,7 @@ import (
// mustOpenView returns a new instance of View with a temporary path.
func mustOpenView(index, field, name string) *view {
path, err := ioutil.TempDir("", "pilosa-view-")
path, err := ioutil.TempDir(*TempDir, "pilosa-view-")
if err != nil {
panic(err)
}