featurebase/server/server_test.go
Matt Jaffee 1c1204fc77
modify PQL parser to handle escapes in string values
This modifies the parser to properly "unquote" incoming strings. So if
a string comes in double or single quoted, we approximately follow Go
rules for removing the quotes and processing escape sequences.

The differences from Go are:
1. we only support backslash, quote, tab and newline escape
sequenences.
2. Single quoted strings are supported and work just like double
quoted strings.
3. The peg parser won't actually accept backquoted strings (I don't
think)

Fixes: #411
2020-05-29 07:58:50 -05:00

1234 lines
37 KiB
Go

// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package server_test
import (
"bytes"
"context"
"encoding/json"
"flag"
"fmt"
"io/ioutil"
"math/rand"
"reflect"
"sort"
"strconv"
"strings"
"testing"
"time"
"github.com/pelletier/go-toml"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/http"
"github.com/pilosa/pilosa/v2/pql"
"github.com/pilosa/pilosa/v2/roaring"
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/test"
"github.com/pkg/errors"
"golang.org/x/sync/errgroup"
)
var runStress bool
func init() { // nolint: gochecknoinits
flag.BoolVar(&runStress, "stress", false, "Enable stress tests (time consuming)")
}
// Ensure program can process queries and maintain consistency.
func TestMain_Set_Quick(t *testing.T) {
if testing.Short() {
t.Skip("short")
}
for i := 0; i < 100; i++ {
t.Run(fmt.Sprint(i), func(t *testing.T) {
t.Parallel()
rand := rand.New(rand.NewSource(int64(i)))
cmds := GenerateSetCommands(1000, rand)
m := test.RunCommand(t)
defer m.Close()
// Create client.
client, err := http.NewInternalClient(m.API.Node().URI.HostPort(), http.GetHTTPClient(nil))
if err != nil {
t.Fatal(err)
}
// Execute Set() commands.
for _, cmd := range cmds {
if err := client.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
t.Fatal(err)
}
if err := client.CreateField(context.Background(), "i", cmd.Field); err != nil && err != pilosa.ErrFieldExists {
t.Fatal(err)
}
if _, err := m.Query(t, "i", "", fmt.Sprintf(`Set(%d, %s=%d)`, cmd.ColumnID, cmd.Field, cmd.ID)); err != nil {
t.Fatal(err)
}
}
// Validate data.
for field, fieldSet := range SetCommands(cmds).Fields() {
for id, columnIDs := range fieldSet {
exp := MustMarshalJSON(map[string]interface{}{
"results": []interface{}{
map[string]interface{}{
"columns": columnIDs,
"attrs": map[string]interface{}{},
},
},
}) + "\n"
if res, err := m.Query(t, "i", "", fmt.Sprintf(`Row(%s=%d)`, field, id)); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result:\n\ngot=%s\n\nexp=%s\n\n", res, exp)
}
}
}
if err := m.Reopen(); err != nil {
t.Fatal(err)
}
// Validate data after reopening.
for field, fieldSet := range SetCommands(cmds).Fields() {
for id, columnIDs := range fieldSet {
exp := MustMarshalJSON(map[string]interface{}{
"results": []interface{}{
map[string]interface{}{
"columns": columnIDs,
"attrs": map[string]interface{}{},
},
},
}) + "\n"
if res, err := m.Query(t, "i", "", fmt.Sprintf(`Row(%s=%d)`, field, id)); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result (reopen):\n\ngot=%s\n\nexp=%s\n\n", res, exp)
}
}
}
})
}
}
// Ensure program can set row attributes and retrieve them.
func TestMain_SetRowAttrs(t *testing.T) {
m := test.RunCommand(t)
defer m.Close()
// Create fields.
client := m.Client()
if err := client.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
t.Fatal(err)
} else if err := client.CreateField(context.Background(), "i", "x"); err != nil {
t.Fatal(err)
} else if err := client.CreateField(context.Background(), "i", "z"); err != nil {
t.Fatal(err)
} else if err := client.CreateField(context.Background(), "i", "neg"); err != nil {
t.Fatal(err)
}
// Set columns on different rows in different fields.
if _, err := m.Query(t, "i", "", `Set(100, x=1)`); err != nil {
t.Fatal(err)
} else if _, err := m.Query(t, "i", "", `Set(100, x=2)`); err != nil {
t.Fatal(err)
} else if _, err := m.Query(t, "i", "", `Set(100, x=2)`); err != nil {
t.Fatal(err)
} else if _, err := m.Query(t, "i", "", `Set(100, neg=3)`); err != nil {
t.Fatal(err)
}
// Set row attributes.
if _, err := m.Query(t, "i", "", `SetRowAttrs(x, 1, x=100)`); err != nil {
t.Fatal(err)
} else if _, err := m.Query(t, "i", "", `SetRowAttrs(x, 2, x=-200)`); err != nil {
t.Fatal(err)
} else if _, err := m.Query(t, "i", "", `SetRowAttrs(z, 2, x=300)`); err != nil {
t.Fatal(err)
} else if _, err := m.Query(t, "i", "", `SetRowAttrs(neg, 3, x=-0.44)`); err != nil {
t.Fatal(err)
}
// Query row x/1.
if res, err := m.Query(t, "i", "", `Row(x=1)`); err != nil {
t.Fatal(err)
} else if res != `{"results":[{"attrs":{"x":100},"columns":[100]}]}`+"\n" {
t.Fatalf("unexpected result: %s", res)
}
// Query row x/2.
if res, err := m.Query(t, "i", "", `Row(x=2)`); err != nil {
t.Fatal(err)
} else if res != `{"results":[{"attrs":{"x":-200},"columns":[100]}]}`+"\n" {
t.Fatalf("unexpected result: %s", res)
}
if err := m.Reopen(); err != nil {
t.Fatal(err)
}
// Query rows after reopening.
if res, err := m.Query(t, "i", "columnAttrs=true", `Row(x=1)`); err != nil {
t.Fatal(err)
} else if res != `{"results":[{"attrs":{"x":100},"columns":[100]}]}`+"\n" {
t.Fatalf("unexpected result(reopen): %s", res)
}
if res, err := m.Query(t, "i", "columnAttrs=true", `Row(neg=3)`); err != nil {
t.Fatal(err)
} else if res != `{"results":[{"attrs":{"x":-0.44},"columns":[100]}]}`+"\n" {
t.Fatalf("unexpected result(reopen): %s", res)
}
// Query row x/2.
if res, err := m.Query(t, "i", "", `Row(x=2)`); err != nil {
t.Fatal(err)
} else if res != `{"results":[{"attrs":{"x":-200},"columns":[100]}]}`+"\n" {
t.Fatalf("unexpected result: %s", res)
}
}
// Ensure program can set column attributes and retrieve them.
func TestMain_SetColumnAttrs(t *testing.T) {
m := test.RunCommand(t)
defer m.Close()
// Create fields.
client := m.Client()
if err := client.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
t.Fatal(err)
} else if err := client.CreateField(context.Background(), "i", "x"); err != nil {
t.Fatal(err)
}
// Set columns on row.
if _, err := m.Query(t, "i", "", `Set(100, x=1)`); err != nil {
t.Fatal(err)
} else if _, err := m.Query(t, "i", "", `Set(101, x=1)`); err != nil {
t.Fatal(err)
}
// Set column attributes.
if _, err := m.Query(t, "i", "", `SetColumnAttrs(100, foo="bar")`); err != nil {
t.Fatal(err)
}
// Query row.
if res, err := m.Query(t, "i", "columnAttrs=true", `Row(x=1)`); err != nil {
t.Fatal(err)
} else if res != `{"results":[{"attrs":{},"columns":[100,101]}],"columnAttrs":[{"id":100,"attrs":{"foo":"bar"}}]}`+"\n" {
t.Fatalf("unexpected result: %s", res)
}
if err := m.Reopen(); err != nil {
t.Fatal(err)
}
// Query row after reopening.
if res, err := m.Query(t, "i", "columnAttrs=true", `Row(x=1)`); err != nil {
t.Fatal(err)
} else if res != `{"results":[{"attrs":{},"columns":[100,101]}],"columnAttrs":[{"id":100,"attrs":{"foo":"bar"}}]}`+"\n" {
t.Fatalf("unexpected result(reopen): %s", res)
}
}
func TestMain_GroupBy(t *testing.T) {
m := test.RunCommand(t)
defer m.Close()
// Create fields.
client := m.Client()
if err := client.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
t.Fatal(err)
}
if err := client.CreateFieldWithOptions(context.Background(), "i", "generalk", pilosa.FieldOptions{Keys: true}); err != nil {
t.Fatal(err)
}
if err := client.CreateFieldWithOptions(context.Background(), "i", "subk", pilosa.FieldOptions{Keys: true}); err != nil {
t.Fatal(err)
}
query := `
Set(0, generalk="ten")
Set(1, generalk="ten")
Set(1001, generalk="ten")
Set(2, generalk="eleven")
Set(1002, generalk="eleven")
Set(2, generalk="twelve")
Set(1002, generalk="twelve")
Set(0, subk="one-hundred")
Set(1, subk="one-hundred")
Set(3, subk="one-hundred")
Set(1001, subk="one-hundred")
Set(2, subk="one-hundred-ten")
Set(0, subk="one-hundred-ten")
`
// Set columns on row.
if _, err := m.Query(t, "i", "", query); err != nil {
t.Fatal(err)
}
expected := []pilosa.GroupCount{
{Group: []pilosa.FieldRow{{Field: "generalk", RowKey: "ten"}, {Field: "subk", RowKey: "one-hundred"}}, Count: 3},
{Group: []pilosa.FieldRow{{Field: "generalk", RowKey: "ten"}, {Field: "subk", RowKey: "one-hundred-ten"}}, Count: 1},
{Group: []pilosa.FieldRow{{Field: "generalk", RowKey: "eleven"}, {Field: "subk", RowKey: "one-hundred-ten"}}, Count: 1},
{Group: []pilosa.FieldRow{{Field: "generalk", RowKey: "twelve"}, {Field: "subk", RowKey: "one-hundred-ten"}}, Count: 1},
}
// Query row.
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))
}
}
func TestMain_MinMaxFloat(t *testing.T) {
m := test.RunCommand(t)
defer m.Close()
// Create fields.
client := m.Client()
if err := client.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
t.Fatal(err)
}
if err := client.CreateFieldWithOptions(context.Background(), "i", "dec", pilosa.FieldOptions{Type: pilosa.FieldTypeDecimal, Scale: 3, Max: pql.NewDecimal(100000, 0)}); err != nil {
t.Fatal(err)
}
query := `
Set(0, dec=1.32)
Set(1, dec=4.44)
`
// Set columns on row.
if _, err := m.Query(t, "i", "", query); err != nil {
t.Fatal(err)
}
// Query row.
exp0 := pilosa.ValCount{DecimalVal: &pql.Decimal{Value: 4440, Scale: 3}, Count: 1}
exp1 := pilosa.ValCount{DecimalVal: &pql.Decimal{Value: 1320, Scale: 3}, Count: 1}
if res, err := m.QueryProtobuf("i", `Max(field=dec) Min(field=dec)`); err != nil {
t.Fatal(err)
} else if !reflect.DeepEqual(res.Results[0], exp0) || !reflect.DeepEqual(res.Results[1], exp1) {
t.Fatalf("unexpected results: %+v", res.Results)
}
}
// Ensure the host can be parsed.
func TestConfig_Parse_Host(t *testing.T) {
if c, err := ParseConfig(`bind = "local"`); err != nil {
t.Fatal(err)
} else if c.Bind != "local" {
t.Fatalf("unexpected host: %s", c.Bind)
}
}
// Ensure the data directory can be parsed.
func TestConfig_Parse_DataDir(t *testing.T) {
if c, err := ParseConfig(`data-dir = "/tmp/foo"`); err != nil {
t.Fatal(err)
} else if c.DataDir != "/tmp/foo" {
t.Fatalf("unexpected data dir: %s", c.DataDir)
}
}
func TestConcurrentFieldCreation(t *testing.T) {
cluster := test.MustRunCluster(t, 3)
defer cluster.Close()
api0 := cluster[0].API
if _, err := api0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil {
t.Fatalf("creating index: %v", err)
}
eg := errgroup.Group{}
for i := 0; i < 100; i++ {
i := i
eg.Go(func() error {
if _, err := api0.CreateField(context.Background(), "i", fmt.Sprintf("f%d", i)); err != nil {
return err
}
return nil
})
}
err := eg.Wait()
if err != nil {
t.Fatalf("creating concurrent field: %v", err)
}
}
func TestTransactionsAPI(t *testing.T) {
cluster := test.MustRunCluster(t, 3)
defer cluster.Close()
api0 := cluster[0].API
api1 := cluster[1].API
ctx := context.Background()
//api2 := cluster[2].API
// can fetch empty transactions
if trnsMap, err := api0.Transactions(ctx); err != nil {
t.Fatalf("getting transactions: %v", err)
} else if len(trnsMap) != 0 {
t.Fatalf("unexpectedly has transactions: %v", trnsMap)
}
// can't fetch transactions from non-coordinator
if _, err := api1.Transactions(ctx); err != pilosa.ErrNodeNotCoordinator {
t.Errorf("api1 should return ErrNodeNotCoordinator when asked for transactions but got: %v", err)
}
// can start transaction
if trns, err := api0.StartTransaction(ctx, "a", time.Minute, false, false); err != nil {
t.Errorf("couldn't start transaction: %v", err)
} else {
test.CompareTransactions(t, &pilosa.Transaction{ID: "a", Active: true, Timeout: time.Minute, Deadline: time.Now().Add(time.Minute)}, trns)
}
// can retrieve transaction from other nodes with remote=true
if trns, err := api1.GetTransaction(ctx, "a", true); err != nil {
t.Errorf("couldn't fetch transaction from other node with remote=true: %v", err)
} else {
test.CompareTransactions(t, &pilosa.Transaction{ID: "a", Active: true, Timeout: time.Minute, Deadline: time.Now().Add(time.Minute)}, trns)
}
// can start transaction with blank id and get uuid back
id := ""
if trns, err := api0.StartTransaction(ctx, id, time.Minute, false, false); err != nil {
t.Errorf("couldn't start transaction: %v", err)
} else {
id = trns.ID
if len(id) != 36 { // UUID
t.Errorf("unexpected generated ID: %s", id)
}
test.CompareTransactions(t, &pilosa.Transaction{ID: id, Active: true, Timeout: time.Minute, Deadline: time.Now().Add(time.Minute)}, trns)
}
// can't finish transaction on non-coordinator
if _, err := api1.FinishTransaction(ctx, id, false); err != pilosa.ErrNodeNotCoordinator {
t.Errorf("unexpected error is not ErrNodeNotCoordinator: %v", err)
}
// can finish transaction
if _, err := api0.FinishTransaction(ctx, id, false); err != nil {
t.Errorf("couldn't finish transaction: %v", err)
}
// can finish previous transaction
if _, err := api0.FinishTransaction(ctx, "a", false); err != nil {
t.Errorf("couldn't finish transaction a: %v", err)
}
// can start exclusive transaction
if te, err := api0.StartTransaction(ctx, "exc", time.Minute, true, false); err != nil {
t.Errorf("couldn't start exclusive transaction: %v", err)
} else if !te.Active {
t.Errorf("expected exclusive transaction to be active: %+v", te)
}
// can finish exclusive transaction
if _, err := api0.FinishTransaction(ctx, "exc", false); err != nil {
t.Errorf("couldn't finish exclusive transaction: %v", err)
}
// can start transaction (with same name as previous finished transaction)
if trns, err := api0.StartTransaction(ctx, "a", time.Minute, false, false); err != nil {
t.Errorf("couldn't start transaction: %v", err)
} else {
test.CompareTransactions(t, &pilosa.Transaction{ID: "a", Active: true, Timeout: time.Minute, Deadline: time.Now().Add(time.Minute)}, trns)
}
// can start exclusive transaction and is not immediately active
if te, err := api0.StartTransaction(ctx, "exc", time.Minute, true, false); err != nil {
t.Errorf("couldn't start exclusive transaction: %v", err)
} else if te.Active {
t.Errorf("expected exclusive transaction to be inactive: %+v", te)
}
// can finish non-exclusive transaction
if _, err := api0.FinishTransaction(ctx, "a", false); err != nil {
t.Errorf("couldn't finish transaction a: %v", err)
}
// can poll exclusive transaction and is active
var excTrns *pilosa.Transaction
if trns, err := api0.GetTransaction(ctx, "exc", false); err != nil {
t.Errorf("couldn't poll exclusive transaction: %v", err)
} else {
excTrns = &pilosa.Transaction{ID: "exc", Active: true, Exclusive: true, Timeout: time.Minute, Deadline: time.Now().Add(time.Minute)}
test.CompareTransactions(t, excTrns, trns)
}
// can't start another exclusive transaction
if trns, err := api0.StartTransaction(ctx, "exc2", time.Minute, true, false); errors.Cause(err) != pilosa.ErrTransactionExclusive {
t.Errorf("unexpected error: %v", err)
} else {
// returned transaction should be the exclusive one which is blocking this one
test.CompareTransactions(t, excTrns, trns)
}
// can't keep the second exclusive name but make it nonexclusive and start a transaction
if trns, err := api0.StartTransaction(ctx, "exc2", time.Minute, false, false); errors.Cause(err) != pilosa.ErrTransactionExclusive {
t.Errorf("unexpected error: %v", err)
} else {
test.CompareTransactions(t, excTrns, trns)
}
// transaction is active on other nodes with remote=true
if trns, err := api1.GetTransaction(ctx, "exc", true); err != nil {
t.Errorf("couldn't poll exclusive transaction: %v", err)
} else {
test.CompareTransactions(t, &pilosa.Transaction{ID: "exc", Active: true, Exclusive: true, Timeout: time.Minute, Deadline: time.Now().Add(time.Minute)}, trns)
}
// LATER, test deadline extension on non-coordinator blocks active, exclusive transaction being returned
}
func TestMain_RecalculateHashes(t *testing.T) {
const clusterSize = 5
cluster := test.MustRunCluster(t, clusterSize)
defer cluster.Close()
// Create the schema.
client0 := cluster[0].Client()
if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
t.Fatal("create index:", err)
}
if err := client0.CreateField(context.Background(), "i", "f"); err != nil {
t.Fatal("create field:", err)
}
// Set some columns
data := []string{}
for rowID := 1; rowID < 10; rowID++ {
for columnID := 1; columnID < 100; columnID++ {
data = append(data, fmt.Sprintf(`Set(%d, f=%d)`, columnID, rowID))
}
}
if _, err := cluster[0].Query(t, "i", "", strings.Join(data, "")); err != nil {
t.Fatal("setting columns:", err)
}
// Calculate caches on the first node
err := cluster[0].RecalculateCaches(t)
if err != nil {
t.Fatalf("recalculating caches: %v", err)
}
target := `{"results":[[{"id":7,"key":"","count":99},{"id":1,"key":"","count":99},{"id":9,"key":"","count":99},{"id":5,"key":"","count":99},{"id":4,"key":"","count":99},{"id":8,"key":"","count":99},{"id":2,"key":"","count":99},{"id":6,"key":"","count":99},{"id":3,"key":"","count":99}]]}`
// Run a TopN query on all nodes. The result should be the same as the target.
for _, m := range cluster {
res, err := m.Query(t, "i", "", `TopN(f)`)
if err != nil {
t.Fatal(err)
}
res = strings.TrimSpace(res)
if sortedString(target) != sortedString(res) {
t.Fatalf("%v != %v", target, res)
}
}
}
// SetCommand represents a command to set a column.
type SetCommand struct {
ID uint64
Field string
ColumnID uint64
}
type SetCommands []SetCommand
// Fields returns the set of column ids for each field/row.
func (a SetCommands) Fields() map[string]map[uint64][]uint64 {
// Create a set of unique commands.
m := make(map[SetCommand]struct{})
for _, cmd := range a {
m[cmd] = struct{}{}
}
// Build unique ids for each field & row.
fields := make(map[string]map[uint64][]uint64)
for cmd := range m {
if fields[cmd.Field] == nil {
fields[cmd.Field] = make(map[uint64][]uint64)
}
fields[cmd.Field][cmd.ID] = append(fields[cmd.Field][cmd.ID], cmd.ColumnID)
}
// Sort each set of column ids.
for _, field := range fields {
for id := range field {
sort.Sort(uint64Slice(field[id]))
}
}
return fields
}
// GenerateSetCommands generates random SetCommand objects.
func GenerateSetCommands(n int, rand *rand.Rand) []SetCommand {
cmds := make([]SetCommand, rand.Intn(n))
for i := range cmds {
cmds[i] = SetCommand{
ID: uint64(rand.Intn(1000)),
Field: "x",
ColumnID: uint64(rand.Intn(10)),
}
}
return cmds
}
// ParseConfig parses s into a Config.
func ParseConfig(s string) (server.Config, error) {
var c server.Config
err := toml.Unmarshal([]byte(s), &c)
return c, err
}
// MustMarshalJSON marshals v into a string. Panic on error.
func MustMarshalJSON(v interface{}) string {
buf, err := json.Marshal(v)
if err != nil {
panic(err)
}
return string(buf)
}
func sortedString(s string) string {
arr := strings.Split(s, "")
sort.Strings(arr)
return strings.Join(arr, "")
}
// uint64Slice represents a sortable slice of uint64 numbers.
type uint64Slice []uint64
func (p uint64Slice) Swap(i, j int) { p[i], p[j] = p[j], p[i] }
func (p uint64Slice) Len() int { return len(p) }
func (p uint64Slice) Less(i, j int) bool { return p[i] < p[j] }
func TestClusteringNodesReplica1(t *testing.T) {
cluster := test.MustRunCluster(t, 3)
defer cluster.Close()
var wait = true
for wait {
wait = false
for _, node := range cluster {
if node.API.State() != pilosa.ClusterStateNormal {
wait = true
}
}
time.Sleep(time.Millisecond * 1)
}
if err := cluster[2].Command.Close(); err != nil {
t.Fatalf("closing third node: %v", err)
}
// confirm that cluster stops accepting queries after one node closes
if _, err := cluster[0].API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") {
t.Fatalf("got unexpected error querying an incomplete cluster: %v", err)
}
// Create new main with the same config.
config := cluster[2].Command.Config
config.Translation.MapSize = 100000
// this isn't necessary, but makes the test run way faster
config.Gossip.Port = strconv.Itoa(int(cluster[2].Command.GossipTransport().URI.Port))
cluster[2].Command = server.NewCommand(cluster[2].Stdin, cluster[2].Stdout, cluster[2].Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore)))
cluster[2].Command.Config = config
// Run new program.
if err := cluster[2].Start(); err != nil {
t.Fatalf("restarting node 2: %v", err)
}
for wait {
wait = false
for _, node := range cluster {
if node.API.State() != pilosa.ClusterStateNormal {
wait = true
}
}
time.Sleep(time.Millisecond)
}
}
func TestClusteringNodesReplica2(t *testing.T) {
cluster := test.MustNewCluster(t, 3)
for _, c := range cluster {
c.Config.Cluster.ReplicaN = 2
}
err := cluster.Start()
if err != nil {
t.Fatalf("starting cluster: %v", err)
}
var wait = true
for wait {
wait = false
for _, node := range cluster {
if node.API.State() != pilosa.ClusterStateNormal {
wait = true
}
}
time.Sleep(time.Millisecond * 1)
}
if err := cluster[2].Command.Close(); err != nil {
t.Fatalf("closing third node: %v", err)
}
if cluster[0].API.State() != pilosa.ClusterStateDegraded {
t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State())
}
// confirm that cluster keeps accepting queries if replication > 1
if _, err := cluster[0].API.CreateIndex(context.Background(), "anewindex", pilosa.IndexOptions{}); err != nil {
t.Fatalf("got unexpected error creating index: %v", err)
}
// confirm that cluster stops accepting queries if 2 nodes fail and replication == 2
if err := cluster[1].Command.Close(); err != nil {
t.Fatalf("closing 2nd node: %v", err)
}
if cluster[0].API.State() != pilosa.ClusterStateStarting {
t.Fatalf("expected state to be Starting, but got %s", cluster[0].API.State())
}
if _, err := cluster[0].API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") {
t.Fatalf("got unexpected error querying an incomplete cluster: %v", err)
}
// Create new main with the same config.
config := cluster[2].Command.Config
config.Translation.MapSize = 100000
// config.Bind = cluster[2].API.Node().URI.HostPort()
// this isn't necessary, but makes the test run way faster
config.Gossip.Port = strconv.Itoa(int(cluster[2].Command.GossipTransport().URI.Port))
cluster[2].Command = server.NewCommand(cluster[2].Stdin, cluster[2].Stdout, cluster[2].Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore)))
cluster[2].Command.Config = config
// Run new program.
if err := cluster[2].Start(); err != nil {
t.Fatalf("restarting node 2: %v", err)
}
if cluster[0].API.State() != pilosa.ClusterStateDegraded {
t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State())
}
// Create new main with the same config.
config = cluster[1].Command.Config
// config.Bind = cluster[1].API.Node().URI.HostPort()
config.Translation.MapSize = 100000
// this isn't necessary, but makes the test run way faster
config.Gossip.Port = strconv.Itoa(int(cluster[1].Command.GossipTransport().URI.Port))
cluster[1].Command = server.NewCommand(cluster[1].Stdin, cluster[1].Stdout, cluster[1].Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore)))
cluster[1].Command.Config = config
// Run new program.
if err := cluster[1].Start(); err != nil {
t.Fatalf("restarting node 2: %v", err)
}
defer cluster.Close()
for wait {
wait = false
for _, node := range cluster {
if node.API.State() != pilosa.ClusterStateNormal {
wait = true
}
}
time.Sleep(time.Millisecond)
}
}
func TestRemoveNodeAfterItDies(t *testing.T) {
cluster := test.MustNewCluster(t, 3)
for _, c := range cluster {
c.Config.Cluster.ReplicaN = 2
}
err := cluster.Start()
if err != nil {
t.Fatalf("starting cluster: %v", err)
}
var wait = true
for wait {
wait = false
for _, node := range cluster {
if node.API.State() != pilosa.ClusterStateNormal {
wait = true
}
}
time.Sleep(time.Millisecond * 1)
}
if err := cluster[2].Command.Close(); err != nil {
t.Fatalf("closing third node: %v", err)
}
if cluster[0].API.State() != pilosa.ClusterStateDegraded {
t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State())
}
if _, err := cluster[0].API.RemoveNode(cluster[2].API.Node().ID); err != nil {
t.Fatalf("removing failed node: %v", err)
}
if cluster[0].API.State() != pilosa.ClusterStateNormal {
t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State())
}
hosts := cluster[0].API.Hosts(context.Background())
if len(hosts) != 2 {
t.Fatalf("unexpected hosts: %v", hosts)
}
}
func TestRemoveConcurrentIndexCreation(t *testing.T) {
cluster := test.MustNewCluster(t, 3)
for _, c := range cluster {
c.Config.Cluster.ReplicaN = 2
}
err := cluster.Start()
if err != nil {
t.Fatalf("starting cluster: %v", err)
}
var wait = true
for wait {
wait = false
for _, node := range cluster {
if node.API.State() != pilosa.ClusterStateNormal {
wait = true
}
}
time.Sleep(time.Millisecond * 1)
}
errc := make(chan error)
go func() {
_, err := cluster[0].API.CreateIndex(context.Background(), "blah", pilosa.IndexOptions{})
errc <- err
}()
if _, err := cluster[0].API.RemoveNode(cluster[2].API.Node().ID); err != nil {
t.Fatalf("removing node: %v", err)
}
for i := 0; cluster[0].API.State() != pilosa.ClusterStateNormal; i++ {
time.Sleep(time.Millisecond)
if i > 10 {
t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State())
}
}
hosts := cluster[0].API.Hosts(context.Background())
if len(hosts) != 2 {
t.Fatalf("unexpected hosts: %v", hosts)
}
if err := <-errc; err != nil {
t.Fatalf("error from index creation: %v", err)
}
}
// Ensure program imports timestamps as UTC.
func TestMain_ImportTimestamp(t *testing.T) {
m := test.RunCommand(t)
defer m.Close()
indexName := "i"
fieldName := "f"
// Create index.
if _, err := m.API.CreateIndex(context.Background(), indexName, pilosa.IndexOptions{}); err != nil {
t.Fatal(err)
}
// Create field.
if _, err := m.API.CreateField(context.Background(), indexName, fieldName, pilosa.OptFieldTypeTime(pilosa.TimeQuantum("YMD"))); err != nil {
t.Fatal(err)
}
data := pilosa.ImportRequest{
Index: indexName,
Field: fieldName,
Shard: 0,
RowIDs: []uint64{1, 2},
ColumnIDs: []uint64{1, 2},
Timestamps: []int64{1514764800000000000, 1577833200000000000}, // 2018-01-01T00:00, 2019-12-31T23:00
}
// Import data.
if err := m.API.Import(context.Background(), &data); err != nil {
t.Fatal(err)
}
// Ensure the correct views were created.
dir := fmt.Sprintf("%s/%s/%s/views", m.Config.DataDir, indexName, fieldName)
files, err := ioutil.ReadDir(dir)
if err != nil {
t.Fatal(err)
}
exp := []string{
"standard", "standard_2018", "standard_201801", "standard_20180101",
"standard_2019", "standard_201912", "standard_20191231",
}
got := []string{}
for _, f := range files {
got = append(got, f.Name())
}
if !reflect.DeepEqual(got, exp) {
t.Fatalf("expected %v, but got %v", exp, got)
}
}
func TestMain_ImportTimestampNoStandardView(t *testing.T) {
m := test.RunCommand(t)
defer m.Close()
indexName := "i"
fieldName := "f-no-standard"
// Create index.
if _, err := m.API.CreateIndex(context.Background(), indexName, pilosa.IndexOptions{}); err != nil {
t.Fatal(err)
}
// Create field.
if _, err := m.API.CreateField(context.Background(), indexName, fieldName, pilosa.OptFieldTypeTime(pilosa.TimeQuantum("YMD"), true)); err != nil {
t.Fatal(err)
}
data := pilosa.ImportRequest{
Index: indexName,
Field: fieldName,
Shard: 0,
RowIDs: []uint64{1, 2},
ColumnIDs: []uint64{1, 2},
Timestamps: []int64{1514764800000000000, 1577833200000000000}, // 2018-01-01T00:00, 2019-12-31T23:00
}
// Import data.
if err := m.API.Import(context.Background(), &data); err != nil {
t.Fatal(err)
}
// Ensure the correct views were created.
dir := fmt.Sprintf("%s/%s/%s/views", m.Config.DataDir, indexName, fieldName)
files, err := ioutil.ReadDir(dir)
if err != nil {
t.Fatal(err)
}
exp := []string{
"standard_2018", "standard_201801", "standard_20180101",
"standard_2019", "standard_201912", "standard_20191231",
}
got := []string{}
for _, f := range files {
got = append(got, f.Name())
}
if !reflect.DeepEqual(got, exp) {
t.Fatalf("expected %v, but got %v", exp, got)
}
}
func TestClusterQueriesAfterRestart(t *testing.T) {
cluster := test.MustRunCluster(t, 3)
defer cluster.Close()
cmd1 := cluster[1]
for _, com := range cluster {
nodes := com.API.Hosts(context.Background())
for _, n := range nodes {
if n.State != "READY" {
t.Fatalf("unexpected node state after upping cluster: %v", nodes)
}
}
}
cmd1.MustCreateIndex(t, "testidx", pilosa.IndexOptions{})
cmd1.MustCreateField(t, "testidx", "testfield", pilosa.OptFieldTypeSet(pilosa.CacheTypeRanked, 10))
// build a query to set the first bit in 100 shards
query := strings.Builder{}
for i := 0; i < 100; i++ {
query.WriteString(fmt.Sprintf("Set(%d, testfield=0)", i*pilosa.ShardWidth))
}
_, err := cmd1.API.Query(context.Background(), &pilosa.QueryRequest{
Index: "testidx",
Query: query.String(),
})
if err != nil {
t.Fatalf("setting 100 bits in 100 shards: %v", err)
}
results, err := cmd1.API.Query(context.Background(), &pilosa.QueryRequest{
Index: "testidx",
Query: "Count(Row(testfield=0))",
})
if err != nil {
t.Fatalf("counting row: %v", err)
}
if results.Results[0].(uint64) != 100 {
t.Fatalf("Count should be 100, but got %v of type %[1]T", results.Results[0])
}
err = cmd1.Command.Close()
if err != nil {
t.Fatalf("closing node0: %v", err)
}
// confirm that cluster stops accepting queries after one node closes
if _, err := cluster[0].API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") {
t.Fatalf("got unexpected error querying an incomplete cluster: %v", err)
}
// Create new main with the same config.
config := cmd1.Command.Config
config.Bind = cmd1.API.Node().URI.HostPort()
// this isn't necessary, but makes the test run way faster
config.Gossip.Port = strconv.Itoa(int(cmd1.Command.GossipTransport().URI.Port))
cmd1.Command = server.NewCommand(cmd1.Stdin, cmd1.Stdout, cmd1.Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore)))
cmd1.Command.Config = config
err = cmd1.Start()
if err != nil {
t.Fatalf("reopening node 0: %v", err)
}
for cmd1.API.State() != pilosa.ClusterStateNormal {
time.Sleep(time.Millisecond)
}
results, err = cmd1.API.Query(context.Background(), &pilosa.QueryRequest{
Index: "testidx",
Query: "Count(Row(testfield=0))",
})
if err != nil {
t.Fatalf("counting row: %v", err)
}
if results.Results[0].(uint64) != 100 {
t.Fatalf("Count should be 100, but got %v of type %[1]T", results.Results[0])
}
}
// TODO: confirm that things keep working if a node is hard-closed (no nodeLeave event) and immediately restarted with a different address.
func TestClusterExhaustingConnections(t *testing.T) {
if !runStress {
t.Skip("stress")
}
cluster := test.MustRunCluster(t, 5)
defer cluster.Close()
cmd1 := cluster[1]
for _, com := range cluster {
nodes := com.API.Hosts(context.Background())
for _, n := range nodes {
if n.State != "READY" {
t.Fatalf("unexpected node state after upping cluster: %v", nodes)
}
}
}
cmd1.MustCreateIndex(t, "testidx", pilosa.IndexOptions{})
cmd1.MustCreateField(t, "testidx", "testfield", pilosa.OptFieldTypeSet(pilosa.CacheTypeRanked, 10))
eg := errgroup.Group{}
for i := 0; i < 20; i++ {
i := i
eg.Go(func() error {
for j := i; j < 10000; j += 20 {
_, err := cluster[i%5].API.Query(context.Background(), &pilosa.QueryRequest{
Index: "testidx",
Query: fmt.Sprintf("Set(%d, testfield=0)", j*pilosa.ShardWidth),
})
if err != nil {
return err
}
}
return nil
})
}
err := eg.Wait()
if err != nil {
t.Fatalf("setting lots of shards: %v", err)
}
}
func TestQueryingWithQuotesAndStuff(t *testing.T) {
m := test.RunCommand(t)
defer m.Close()
client, err := http.NewInternalClient(m.API.Node().URI.HostPort(), http.GetHTTPClient(nil))
if err != nil {
t.Fatal(err)
}
// Execute Set() commands.
if err := client.CreateIndex(context.Background(), "i", pilosa.IndexOptions{Keys: true}); err != nil {
t.Fatal(err)
}
if err := client.CreateFieldWithOptions(context.Background(), "i", "fld", pilosa.FieldOptions{Keys: true}); err != nil {
t.Fatal(err)
}
// Test escaped single quote gets set properly
if res, err := m.Query(t, "i", "", `Set('bl\'ah', fld=ha)`); err != nil {
t.Fatal(err)
} else if !strings.Contains(res, "[true]") {
t.Errorf("setting escaped single quote result: %s", res)
}
if res, err := m.Query(t, "i", "", `Row(fld=ha)`); err != nil {
t.Fatal(err)
} else if !strings.Contains(res, `bl'ah`) {
t.Errorf("value with escaped single quote set improperly: %s", res)
}
// Test escaped double quote gets set properly
if res, err := m.Query(t, "i", "", `Set("d\"ah", fld=dq)`); err != nil {
t.Fatal(err)
} else if !strings.Contains(res, "[true]") {
t.Errorf("value with escaped double quote set improperly: %s", res)
}
if res, err := m.Query(t, "i", "", `Row(fld=dq)`); err != nil {
t.Fatal(err)
} else if !strings.Contains(res, `d\"ah`) {
// the backslash is there because JSON needs to escape the
// double quote since it uses double quotes
t.Errorf("value with escaped double quote set improperly: %s", res)
}
}
func TestClusterExhaustingConnectionsImport(t *testing.T) {
if !runStress {
t.Skip("stress")
}
cluster := test.MustRunCluster(t, 5)
defer cluster.Close()
cmd1 := cluster[1]
for _, com := range cluster {
nodes := com.API.Hosts(context.Background())
for _, n := range nodes {
if n.State != "READY" {
t.Fatalf("unexpected node state after upping cluster: %v", nodes)
}
}
}
cmd1.MustCreateIndex(t, "testidx", pilosa.IndexOptions{})
cmd1.MustCreateField(t, "testidx", "testfield", pilosa.OptFieldTypeSet(pilosa.CacheTypeRanked, 10))
bm := roaring.NewBitmap()
bm.DirectAdd(0)
buf := &bytes.Buffer{}
_, err := bm.WriteTo(buf)
if err != nil {
t.Fatalf("writing to buffer: %v", err)
}
data := buf.Bytes()
eg := errgroup.Group{}
for i := uint64(0); i < 20; i++ {
i := i
eg.Go(func() error {
for j := i; j < 10000; j += 20 {
if (j-i)%1000 == 0 {
fmt.Printf("%d is %.2f%% done.\n", i, float64(j-i)*100/100000)
}
err := cluster[i%5].API.ImportRoaring(context.Background(), "testidx", "testfield", j, false, &pilosa.ImportRoaringRequest{
Views: map[string][]byte{
"": data,
},
})
if err != nil {
return err
}
}
return nil
})
}
err = eg.Wait()
if err != nil {
t.Fatalf("setting lots of shards: %v", err)
}
}
func TestClusterMinMaxSumDecimal(t *testing.T) {
cluster := test.MustRunCluster(t, 3)
defer cluster.Close()
cmd := cluster[0]
cmd.MustCreateIndex(t, "testdec", pilosa.IndexOptions{Keys: true, TrackExistence: true})
cmd.MustCreateField(t, "testdec", "adec", pilosa.OptFieldTypeDecimal(2))
test.Do(t, "POST", cluster[0].URL()+"/index/testdec/query", `
Set("a", adec=42.2)
Set("b", adec=11.12)
Set("c", adec=13.41)
Set("d", adec=99.87)
Set("e", adec=11.13)
Set("f", adec=12.12)
Set("g", adec=15.52)
Set("h", adec=100.22)
`)
result := test.Do(t, "POST", cluster[0].URL()+"/index/testdec/query", "Sum(field=adec)")
if !strings.Contains(result.Body, `"decimalValue":305.59`) {
t.Fatalf("expected decimal sum of 305.59, but got: '%s'", result.Body)
} else if !strings.Contains(result.Body, `"count":8`) {
t.Fatalf("expected count 8, but got: '%s'", result.Body)
}
result = test.Do(t, "POST", cluster[0].URL()+"/index/testdec/query", "Max(field=adec)")
if !strings.Contains(result.Body, `"decimalValue":100.22`) {
t.Fatalf("expected decimal max of 100.22, but got: '%s'", result.Body)
} else if !strings.Contains(result.Body, `"count":1`) {
t.Fatalf("expected count 1, but got: '%s'", result.Body)
}
result = test.Do(t, "POST", cluster[0].URL()+"/index/testdec/query", "Min(field=adec)")
if !strings.Contains(result.Body, `"decimalValue":11.12`) {
t.Fatalf("expected decimal min of 11.12, but got: '%s'", result.Body)
} else if !strings.Contains(result.Body, `"count":1`) {
t.Fatalf("expected count 1, but got: '%s'", result.Body)
}
}