mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-05 08:10:50 +00:00
Merge branch 'master' into seebs/deadlock
This commit is contained in:
commit
3bbd3cbfb1
23 changed files with 1295 additions and 1165 deletions
2
Makefile
2
Makefile
|
|
@ -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 \
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
66
executor.go
66
executor.go
|
|
@ -18,7 +18,6 @@ import (
|
|||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log"
|
||||
"sort"
|
||||
"time"
|
||||
|
||||
|
|
@ -806,7 +805,7 @@ func (e *executor) executeTopNShard(ctx context.Context, index string, c *pql.Ca
|
|||
return nil, nil
|
||||
}
|
||||
|
||||
if minThreshold <= 0 {
|
||||
if minThreshold == 0 {
|
||||
minThreshold = defaultMinThreshold
|
||||
}
|
||||
|
||||
|
|
@ -915,6 +914,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)
|
||||
}
|
||||
|
|
@ -1089,6 +1094,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 {
|
||||
|
|
@ -1097,7 +1113,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.
|
||||
|
|
@ -1122,17 +1138,12 @@ 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 {
|
||||
|
|
@ -1182,7 +1193,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.
|
||||
|
|
@ -1576,14 +1587,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)
|
||||
|
|
@ -1702,19 +1713,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.
|
||||
|
|
@ -1739,15 +1750,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)
|
||||
|
|
@ -1764,7 +1775,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.
|
||||
|
|
@ -1783,7 +1794,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
|
||||
|
||||
|
|
@ -2344,7 +2355,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":
|
||||
|
|
@ -2449,7 +2460,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)
|
||||
|
|
@ -2560,7 +2571,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
|
||||
}
|
||||
|
|
@ -2758,11 +2769,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
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
},
|
||||
}
|
||||
|
|
|
|||
135
executor_test.go
135
executor_test.go
|
|
@ -2283,7 +2283,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 {
|
||||
|
|
@ -3038,22 +3038,29 @@ 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)
|
||||
}
|
||||
|
|
@ -3081,35 +3088,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:",
|
||||
},
|
||||
}
|
||||
|
||||
|
|
@ -3123,7 +3122,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())
|
||||
}
|
||||
})
|
||||
}
|
||||
|
|
@ -3169,68 +3168,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{},
|
||||
},
|
||||
}
|
||||
|
|
@ -3282,13 +3293,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},
|
||||
|
|
@ -3297,7 +3322,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)
|
||||
})
|
||||
|
||||
|
|
@ -3307,7 +3332,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)
|
||||
})
|
||||
|
||||
|
|
@ -3317,7 +3342,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)
|
||||
})
|
||||
|
||||
|
|
@ -3326,7 +3351,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)
|
||||
|
||||
})
|
||||
|
|
@ -3347,7 +3372,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)
|
||||
})
|
||||
|
||||
|
|
@ -3375,7 +3400,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},
|
||||
|
|
@ -3385,14 +3410,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},
|
||||
}
|
||||
|
|
@ -3416,7 +3441,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},
|
||||
|
|
@ -3428,7 +3453,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},
|
||||
|
|
@ -3439,7 +3464,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},
|
||||
|
|
@ -3464,7 +3489,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},
|
||||
|
|
@ -3501,11 +3526,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...)
|
||||
}
|
||||
|
|
@ -3548,7 +3573,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)
|
||||
|
||||
})
|
||||
|
|
@ -3597,7 +3622,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))`)
|
||||
}
|
||||
})
|
||||
|
||||
|
|
@ -3605,7 +3630,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)`)
|
||||
}
|
||||
})
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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{}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
2119
pql/pql.peg.go
2119
pql/pql.peg.go
File diff suppressed because it is too large
Load diff
|
|
@ -3868,7 +3868,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
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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))
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue