mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-06 19:07:50 +00:00
Merge branch 'master' into rbf-typos
This commit is contained in:
commit
ddb8c95737
32 changed files with 617 additions and 337 deletions
25
attr_test.go
25
attr_test.go
|
|
@ -15,7 +15,6 @@
|
|||
package pilosa_test
|
||||
|
||||
import (
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"reflect"
|
||||
"runtime"
|
||||
|
|
@ -24,11 +23,12 @@ import (
|
|||
|
||||
"github.com/pilosa/pilosa/v2"
|
||||
"github.com/pilosa/pilosa/v2/boltdb"
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
)
|
||||
|
||||
// Ensure database can set and retrieve column attributes.
|
||||
func TestAttrStore_Attrs(t *testing.T) {
|
||||
s := MustOpenAttrStore()
|
||||
s := MustOpenAttrStore(t)
|
||||
defer s.Close()
|
||||
|
||||
// Set attributes.
|
||||
|
|
@ -57,7 +57,7 @@ func TestAttrStore_Attrs(t *testing.T) {
|
|||
|
||||
// Ensure database returns a non-nil empty map if unset.
|
||||
func TestAttrStore_Attrs_Empty(t *testing.T) {
|
||||
s := MustOpenAttrStore()
|
||||
s := MustOpenAttrStore(t)
|
||||
defer s.Close()
|
||||
|
||||
if m, err := s.Attrs(100); err != nil {
|
||||
|
|
@ -69,7 +69,7 @@ func TestAttrStore_Attrs_Empty(t *testing.T) {
|
|||
|
||||
// Ensure database can unset attributes if explicitly set to nil.
|
||||
func TestAttrStore_Attrs_Unset(t *testing.T) {
|
||||
s := MustOpenAttrStore()
|
||||
s := MustOpenAttrStore(t)
|
||||
defer s.Close()
|
||||
|
||||
// Set attributes.
|
||||
|
|
@ -89,7 +89,7 @@ func TestAttrStore_Attrs_Unset(t *testing.T) {
|
|||
|
||||
// Ensure attribute block checksums can be returned.
|
||||
func TestAttrStore_Blocks(t *testing.T) {
|
||||
s := MustOpenAttrStore()
|
||||
s := MustOpenAttrStore(t)
|
||||
defer s.Close()
|
||||
|
||||
// Set attributes.
|
||||
|
|
@ -135,19 +135,22 @@ type AttrStore struct {
|
|||
}
|
||||
|
||||
// NewAttrStore returns a new instance of AttrStore.
|
||||
func NewAttrStore(string) pilosa.AttrStore {
|
||||
f, err := ioutil.TempFile("", "pilosa-attr-")
|
||||
func NewAttrStore(tb testing.TB) pilosa.AttrStore {
|
||||
f, err := testhook.TempFile(tb, "pilosa-attr-")
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
// Note, even though the file is closed, TempFile will still avoid
|
||||
// creating the same name again if we leave it existing. The boltdb
|
||||
// code may already be deleting this, so the TestHook deletion
|
||||
// may not matter but it's more reliable this way.
|
||||
f.Close()
|
||||
os.Remove(f.Name())
|
||||
|
||||
return &AttrStore{boltdb.NewAttrStore(f.Name())}
|
||||
}
|
||||
|
||||
func BenchmarkAttrStore_Duplicate(b *testing.B) {
|
||||
s := MustOpenAttrStore()
|
||||
s := MustOpenAttrStore(b)
|
||||
defer s.Close()
|
||||
|
||||
// Set attributes.
|
||||
|
|
@ -186,8 +189,8 @@ func BenchmarkAttrStore_Duplicate(b *testing.B) {
|
|||
}
|
||||
|
||||
// MustOpenAttrStore returns a new, opened attribute store at a temporary path. Panic on error.
|
||||
func MustOpenAttrStore() pilosa.AttrStore {
|
||||
s := NewAttrStore("")
|
||||
func MustOpenAttrStore(tb testing.TB) pilosa.AttrStore {
|
||||
s := NewAttrStore(tb)
|
||||
if err := s.Open(); err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -19,7 +19,6 @@ import (
|
|||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
|
|
@ -32,7 +31,6 @@ import (
|
|||
"runtime/pprof"
|
||||
)
|
||||
|
||||
var _ = ioutil.TempFile
|
||||
var _ = pprof.StartCPUProfile
|
||||
|
||||
var (
|
||||
|
|
|
|||
|
|
@ -17,8 +17,6 @@ import (
|
|||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"reflect"
|
||||
"strconv"
|
||||
"testing"
|
||||
|
|
@ -26,13 +24,14 @@ import (
|
|||
|
||||
"github.com/pilosa/pilosa/v2"
|
||||
"github.com/pilosa/pilosa/v2/boltdb"
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
"github.com/pilosa/pilosa/v2/topology"
|
||||
)
|
||||
|
||||
//var vv = pilosa.VV
|
||||
|
||||
func TestTranslateStore_TranslateKey(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
// Ensure initial key translates to first ID for shard
|
||||
|
|
@ -57,7 +56,7 @@ func TestTranslateStore_TranslateKey(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestTranslateStore_TranslateKeys(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
ids, err := s.TranslateKeys([]string{"abc", "abc"}, true)
|
||||
|
|
@ -97,7 +96,7 @@ func TestTranslateStore_TranslateKeys(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestTranslateStore_CreateKeys(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
ids, err := s.CreateKeys("abc", "abc")
|
||||
|
|
@ -137,7 +136,7 @@ func TestTranslateStore_CreateKeys(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestTranslateStore_ReadKey(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
id, err := s.TranslateKey("foo", false)
|
||||
|
|
@ -172,7 +171,7 @@ func TestTranslateStore_ReadKey(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestTranslateStore_ReadKeys(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
ids, err := s.TranslateKeys([]string{"foo", "bar", "baz", "baz", "bar", "foo"}, false)
|
||||
|
|
@ -200,7 +199,7 @@ func TestTranslateStore_ReadKeys(t *testing.T) {
|
|||
}
|
||||
}
|
||||
func TestTranslateStore_TranslateID(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
// Setup initial keys.
|
||||
|
|
@ -239,7 +238,7 @@ func TestTranslateStore_TranslateID(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestTranslateStore_TranslateIDs(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
// Setup initial keys.
|
||||
|
|
@ -293,7 +292,7 @@ func TestTranslateStore_FindKeys(t *testing.T) {
|
|||
for _, c := range cases {
|
||||
c := c
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
var naiveMap map[string]uint64
|
||||
|
|
@ -339,7 +338,7 @@ func TestTranslateStore_FindKeys(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestTranslateStore_MaxID(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
// Generate a bunch of keys.
|
||||
|
|
@ -364,7 +363,7 @@ func TestTranslateStore_MaxID(t *testing.T) {
|
|||
|
||||
func TestTranslateStore_EntryReader(t *testing.T) {
|
||||
t.Run("OK", func(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
// Create multiple new keys.
|
||||
|
|
@ -422,7 +421,7 @@ func TestTranslateStore_EntryReader(t *testing.T) {
|
|||
|
||||
// Ensure reader will read as soon as a new write comes in using WriteNotify().
|
||||
t.Run("WriteNotify", func(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
// Start reader from initial position.
|
||||
|
|
@ -465,7 +464,7 @@ func TestTranslateStore_EntryReader(t *testing.T) {
|
|||
|
||||
// Ensure exits read on close.
|
||||
t.Run("Close", func(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
// Start reader from initial position.
|
||||
|
|
@ -499,7 +498,7 @@ func TestTranslateStore_EntryReader(t *testing.T) {
|
|||
|
||||
// Ensure exits read on store close.
|
||||
t.Run("StoreClose", func(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
// Start reader from initial position.
|
||||
|
|
@ -533,8 +532,8 @@ func TestTranslateStore_EntryReader(t *testing.T) {
|
|||
}
|
||||
|
||||
// MustNewTranslateStore returns a new TranslateStore with a temporary path.
|
||||
func MustNewTranslateStore() *boltdb.TranslateStore {
|
||||
f, err := ioutil.TempFile("", "")
|
||||
func MustNewTranslateStore(tb testing.TB) *boltdb.TranslateStore {
|
||||
f, err := testhook.TempFile(tb, "translate-store")
|
||||
if err != nil {
|
||||
panic(err)
|
||||
} else if err := f.Close(); err != nil {
|
||||
|
|
@ -548,7 +547,7 @@ func MustNewTranslateStore() *boltdb.TranslateStore {
|
|||
|
||||
func TestTranslateStore_ReadWrite(t *testing.T) {
|
||||
t.Run("WriteTo_ReadFrom", func(t *testing.T) {
|
||||
s := MustOpenNewTranslateStore()
|
||||
s := MustOpenNewTranslateStore(t)
|
||||
defer MustCloseTranslateStore(s)
|
||||
|
||||
batch0 := []string{}
|
||||
|
|
@ -612,8 +611,8 @@ func TestTranslateStore_ReadWrite(t *testing.T) {
|
|||
}
|
||||
|
||||
// MustOpenNewTranslateStore returns a new, opened TranslateStore.
|
||||
func MustOpenNewTranslateStore() *boltdb.TranslateStore {
|
||||
s := MustNewTranslateStore()
|
||||
func MustOpenNewTranslateStore(tb testing.TB) *boltdb.TranslateStore {
|
||||
s := MustNewTranslateStore(tb)
|
||||
if err := s.Open(); err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
|
@ -624,7 +623,5 @@ func MustOpenNewTranslateStore() *boltdb.TranslateStore {
|
|||
func MustCloseTranslateStore(s *boltdb.TranslateStore) {
|
||||
if err := s.Close(); err != nil {
|
||||
panic(err)
|
||||
} else if err := os.Remove(s.Path); err != nil {
|
||||
panic(err)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
48
cache.go
48
cache.go
|
|
@ -369,6 +369,12 @@ type Pair struct {
|
|||
Count uint64 `json:"count"`
|
||||
}
|
||||
|
||||
func (pair Pair) MarshalJSON() ([]byte, error) {
|
||||
buffer := new(bytes.Buffer)
|
||||
_, err := buffer.WriteString(fmt.Sprintf(`{"count":%v,"id":%v,"key":"%v"}`, pair.Count, pair.ID, pair.Key))
|
||||
return buffer.Bytes(), err
|
||||
}
|
||||
|
||||
// PairField is a Pair with its associated field.
|
||||
type PairField struct {
|
||||
Pair Pair
|
||||
|
|
@ -561,7 +567,47 @@ func (p *PairsField) ToRows(callback func(*pb.RowResponse) error) error {
|
|||
// MarshalJSON marshals PairsField into a JSON-encoded byte slice,
|
||||
// excluding `Field`.
|
||||
func (p PairsField) MarshalJSON() ([]byte, error) {
|
||||
return json.Marshal(p.Pairs)
|
||||
return json.Marshal(SortedPairs(p.Pairs))
|
||||
}
|
||||
|
||||
type SortedPairs []Pair
|
||||
|
||||
func (sortPair SortedPairs) MarshalJSON() ([]byte, error) {
|
||||
sort.Slice(sortPair, func(i, j int) bool {
|
||||
if sortPair[j].Count == sortPair[i].Count {
|
||||
if sortPair[j].Key != "" {
|
||||
return sortPair[j].Key < sortPair[i].Key
|
||||
}
|
||||
return sortPair[j].ID < sortPair[i].ID
|
||||
}
|
||||
return sortPair[j].Count < sortPair[i].Count
|
||||
})
|
||||
buffer := new(bytes.Buffer)
|
||||
_, err := buffer.Write([]byte("["))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for i, pair := range sortPair {
|
||||
jsonValue, err := json.Marshal(pair)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if i != 0 {
|
||||
_, err = buffer.Write([]byte(","))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
_, err = buffer.Write(jsonValue)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
_, err = buffer.Write([]byte("]"))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return buffer.Bytes(), nil
|
||||
}
|
||||
|
||||
// int64Slice represents a sortable slice of int64 numbers.
|
||||
|
|
|
|||
|
|
@ -212,15 +212,18 @@ NewSetup:
|
|||
pql, err := cfg.GenQuery(index)
|
||||
panicOn(err)
|
||||
|
||||
if cfg.Verbose {
|
||||
fmt.Printf("pql = '%v'\n", pql)
|
||||
}
|
||||
|
||||
// Query node0.
|
||||
res, err := cli.Query(ctx, index, &pilosa.QueryRequest{Index: index, Query: pql})
|
||||
colorReset := "\033[0m"
|
||||
if err != nil {
|
||||
AlwaysPrintf("QUERY FAILED! queries before this=%v; err = '%v', pql='%v'", loops, err, pql)
|
||||
return err
|
||||
colorRed := "\033[31m"
|
||||
colorCyan := "\033[36m"
|
||||
AlwaysPrintf("\n%vDIFF%v queries before this=%v;\nPQL=%v%v%v\nerr='%v'", colorRed, colorReset, loops, colorCyan, pql, colorReset, err)
|
||||
} else {
|
||||
colorYellow := "\033[33m"
|
||||
if cfg.Verbose {
|
||||
fmt.Printf("pql = %v%v%v\n", colorYellow, pql, colorReset)
|
||||
}
|
||||
}
|
||||
if cfg.VeryVerbose {
|
||||
fmt.Printf("success on pql = '%v'; res='%v'\n", pql, res.Results[0])
|
||||
|
|
@ -374,9 +377,27 @@ func (cfg *RandomQueryConfig) Setup(api API) (err error) {
|
|||
}
|
||||
case "int":
|
||||
foundIntField = true
|
||||
fallthrough // I bet you thought you'd never see this used
|
||||
minPQL := fmt.Sprintf("Min(field=%v)", fld.Name)
|
||||
resp, err := api.Query(ctx, ii.Name, &pilosa.QueryRequest{Index: ii.Name, Query: minPQL})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
min := resp.Results[0]
|
||||
maxPQL := fmt.Sprintf("Max(field=%v)", fld.Name)
|
||||
resp, err = api.Query(ctx, ii.Name, &pilosa.QueryRequest{Index: ii.Name, Query: maxPQL})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
max := resp.Results[0]
|
||||
m := pql.Decimal{}
|
||||
x := pql.Decimal{}
|
||||
m.Scale = fld.Options.Scale
|
||||
m.Value = min.(pilosa.ValCount).Val
|
||||
x.Value = max.(pilosa.ValCount).Val
|
||||
x.Scale = fld.Options.Scale
|
||||
|
||||
cfg.AddIntField(ii.Name, fld.Name, m, x, fld.Options.Scale, fld.Options.Type == "decimal")
|
||||
case "decimal":
|
||||
// we'll ignore row keys and just use value ranges
|
||||
cfg.AddIntField(ii.Name, fld.Name, fld.Options.Min, fld.Options.Max, fld.Options.Scale, fld.Options.Type == "decimal")
|
||||
default:
|
||||
AlwaysPrintf("ignoring field %q: unhandled type %q\n", fld.Name, fld.Options.Type)
|
||||
|
|
@ -451,10 +472,6 @@ func (cfg *RandomQueryConfig) AddIntField(index, field string, min, max pql.Deci
|
|||
func (cfg *RandomQueryConfig) GenQuery(index string) (pql string, err error) {
|
||||
|
||||
tree := cfg.GenTree(index, cfg.TreeDepth)
|
||||
|
||||
if cfg.Verbose {
|
||||
fmt.Printf("%v\n", tree.StringIndent(0))
|
||||
}
|
||||
pql = tree.ToPQL()
|
||||
|
||||
// avoid using too much bandwidth, just count the final bitmap.
|
||||
|
|
@ -508,18 +525,16 @@ func (cfg *RandomQueryConfig) GenTree(index string, depth int) (tr *Tree) {
|
|||
tr = &Tree{S: f}
|
||||
numChild := 2
|
||||
switch f {
|
||||
case "Union", "Intersect", "Xor":
|
||||
case "Union", "Intersect":
|
||||
numChild = cfg.Rnd.Intn(8) + 2
|
||||
case "Xor":
|
||||
numChild = cfg.Rnd.Intn(2) + 1
|
||||
case "Not":
|
||||
numChild = 1
|
||||
case "Difference":
|
||||
numChild = 2
|
||||
case "Distinct":
|
||||
numChild = 1
|
||||
// sometimes do a bare distinct without a filter
|
||||
if cfg.Rnd.Intn(10) == 0 {
|
||||
numChild = 0
|
||||
}
|
||||
r = cfg.Rnd.Intn(len(features.Distinctables))
|
||||
tr.Args = append(tr.Args, fmt.Sprintf("field=%s", features.Distinctables[r].Field))
|
||||
}
|
||||
|
|
|
|||
|
|
@ -26,6 +26,7 @@ import (
|
|||
"time"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/cmd"
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
"github.com/spf13/cobra"
|
||||
)
|
||||
|
||||
|
|
@ -135,7 +136,7 @@ func executeDry(t *testing.T, tests []commandTest) {
|
|||
// a temp config file with the cfgFileContent string as its content.
|
||||
func (ct *commandTest) setupCommand(t *testing.T) *cobra.Command {
|
||||
// make config file
|
||||
cfgFile, err := ioutil.TempFile("", "")
|
||||
cfgFile, err := testhook.TempFile(t, "cmdconf")
|
||||
failErr(t, err, "making temp file")
|
||||
_, err = cfgFile.WriteString(ct.cfgFileContent)
|
||||
failErr(t, err, "writing config to temp file")
|
||||
|
|
@ -180,9 +181,9 @@ func TestRootCommand(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestRootCommand_Config(t *testing.T) {
|
||||
file, err := ioutil.TempFile("", "test.conf")
|
||||
file, err := testhook.TempFile(t, "test.conf")
|
||||
if err != nil {
|
||||
panic(err)
|
||||
t.Fatalf("creating config file: %v", err)
|
||||
}
|
||||
config := `data-dir = "/tmp/pil5_0"
|
||||
bind = "127.0.0.1:10101"
|
||||
|
|
|
|||
|
|
@ -16,13 +16,13 @@ package cmd_test
|
|||
|
||||
import (
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/cmd"
|
||||
_ "github.com/pilosa/pilosa/v2/test"
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
"github.com/pilosa/pilosa/v2/toml"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
|
@ -42,9 +42,9 @@ func nextPort() string { //nolint:unused
|
|||
|
||||
func TestServerConfig(t *testing.T) {
|
||||
t.Skip("pilosa hosts config (cmd.Server.Config.Cluster.Hosts and brethren) is test only and will go away with high probability. skip for now.")
|
||||
actualDataDir, err := ioutil.TempDir("", "")
|
||||
actualDataDir, err := testhook.TempDir(t, "")
|
||||
failErr(t, err, "making data dir")
|
||||
logFile, err := ioutil.TempFile("", "")
|
||||
logFile, err := testhook.TempFile(t, "")
|
||||
failErr(t, err, "making log file")
|
||||
tests := []commandTest{
|
||||
// TEST 0
|
||||
|
|
@ -64,7 +64,7 @@ func TestServerConfig(t *testing.T) {
|
|||
bind-grpc = ` + nextPort() + `
|
||||
max-writes-per-request = 3000
|
||||
long-query-time = "1m10s"
|
||||
|
||||
|
||||
[cluster]
|
||||
replicas = 2
|
||||
long-query-time = "1m10s"
|
||||
|
|
@ -189,7 +189,7 @@ func TestServerConfig(t *testing.T) {
|
|||
}
|
||||
func TestServerConfig_DeprecateLongQueryTime(t *testing.T) {
|
||||
t.Skip("pilosa hosts config (cmd.Server.Config.Cluster.Hosts and brethren) is test only and will go away with high probability. skip for now.")
|
||||
actualDataDir, err := ioutil.TempDir("", "")
|
||||
actualDataDir, err := testhook.TempDir(t, "")
|
||||
failErr(t, err, "making data dir")
|
||||
|
||||
tests := []commandTest{
|
||||
|
|
|
|||
|
|
@ -16,20 +16,22 @@ package ctl
|
|||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/hex"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"math/rand"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"context"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
)
|
||||
|
||||
func TestCheckCommand_RunCacheFile(t *testing.T) {
|
||||
cacheFile := TempFileName("test", ".cache")
|
||||
fi, err := testhook.TempFile(t, "test*.cache")
|
||||
if err != nil {
|
||||
t.Fatalf("creating test file: %v", err)
|
||||
}
|
||||
cacheFile := fi.Name()
|
||||
|
||||
rder := []byte{}
|
||||
stdin := bytes.NewReader(rder)
|
||||
|
|
@ -37,7 +39,7 @@ func TestCheckCommand_RunCacheFile(t *testing.T) {
|
|||
cm := NewCheckCommand(stdin, w, w)
|
||||
cm.Paths = []string{cacheFile}
|
||||
|
||||
err := cm.Run(context.Background())
|
||||
err = cm.Run(context.Background())
|
||||
w.Close()
|
||||
var buf bytes.Buffer
|
||||
if _, err := io.Copy(&buf, r); err != nil {
|
||||
|
|
@ -50,7 +52,11 @@ func TestCheckCommand_RunCacheFile(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestCheckCommand_RunSnapshot(t *testing.T) {
|
||||
snapshotFile := TempFileName("test", ".snapshotting")
|
||||
fi, err := testhook.TempFile(t, "test*.snapshotting")
|
||||
if err != nil {
|
||||
t.Fatalf("creating test file: %v", err)
|
||||
}
|
||||
snapshotFile := fi.Name()
|
||||
|
||||
rder := []byte{}
|
||||
stdin := bytes.NewReader(rder)
|
||||
|
|
@ -58,7 +64,7 @@ func TestCheckCommand_RunSnapshot(t *testing.T) {
|
|||
cm := NewCheckCommand(stdin, w, w)
|
||||
cm.Paths = []string{snapshotFile}
|
||||
|
||||
err := cm.Run(context.Background())
|
||||
err = cm.Run(context.Background())
|
||||
w.Close()
|
||||
var buf bytes.Buffer
|
||||
if _, err := io.Copy(&buf, r); err != nil {
|
||||
|
|
@ -71,10 +77,11 @@ func TestCheckCommand_RunSnapshot(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestCheckCommand_Run(t *testing.T) {
|
||||
file, err := ioutil.TempFile("", "")
|
||||
file, err := testhook.TempFile(t, "run-command")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
fname := file.Name()
|
||||
if _, err := file.Write([]byte("1234,1223")); err != nil {
|
||||
t.Fatalf("writing to temp file: %v", err)
|
||||
}
|
||||
|
|
@ -84,7 +91,7 @@ func TestCheckCommand_Run(t *testing.T) {
|
|||
stdin := bytes.NewReader(rder)
|
||||
r, w, _ := os.Pipe()
|
||||
cm := NewCheckCommand(stdin, w, w)
|
||||
cm.Paths = []string{file.Name()}
|
||||
cm.Paths = []string{fname}
|
||||
|
||||
err = cm.Run(context.Background())
|
||||
w.Close()
|
||||
|
|
@ -99,10 +106,3 @@ func TestCheckCommand_Run(t *testing.T) {
|
|||
}
|
||||
// Todo: need correct roaring file for happy path
|
||||
}
|
||||
|
||||
// TempFileName generates a temporary filename with extension
|
||||
func TempFileName(prefix, suffix string) string {
|
||||
randBytes := make([]byte, 16)
|
||||
rand.Read(randBytes)
|
||||
return filepath.Join(os.TempDir(), prefix+hex.EncodeToString(randBytes)+suffix)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -29,6 +29,7 @@ import (
|
|||
|
||||
"github.com/pilosa/pilosa/v2"
|
||||
"github.com/pilosa/pilosa/v2/test"
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
)
|
||||
|
||||
func TestImportCommand_Validation(t *testing.T) {
|
||||
|
|
@ -58,7 +59,7 @@ func TestImportCommand_Basic(t *testing.T) {
|
|||
buf := bytes.Buffer{}
|
||||
stdin, stdout, stderr := GetIO(buf)
|
||||
cm := NewImportCommand(stdin, stdout, stderr)
|
||||
file, err := ioutil.TempFile("", "import.csv")
|
||||
file, err := testhook.TempFile(t, "import.csv")
|
||||
if err != nil {
|
||||
t.Fatalf("creating tempfile: %v", err)
|
||||
}
|
||||
|
|
@ -90,7 +91,7 @@ func TestImportCommand_Basic(t *testing.T) {
|
|||
buf := bytes.Buffer{}
|
||||
stdin, stdout, stderr := GetIO(buf)
|
||||
cm := NewImportCommand(stdin, stdout, stderr)
|
||||
file, err := ioutil.TempFile("", "import.csv")
|
||||
file, err := testhook.TempFile(t, "import.csv")
|
||||
if err != nil {
|
||||
t.Fatalf("creating tempfile: %v", err)
|
||||
}
|
||||
|
|
@ -123,7 +124,7 @@ func TestImportCommand_RunValue(t *testing.T) {
|
|||
buf := bytes.Buffer{}
|
||||
stdin, stdout, stderr := GetIO(buf)
|
||||
cm := NewImportCommand(stdin, stdout, stderr)
|
||||
file, err := ioutil.TempFile("", "import-value.csv")
|
||||
file, err := testhook.TempFile(t, "import-value.csv")
|
||||
if err != nil {
|
||||
t.Fatalf("creating tempfile: %v", err)
|
||||
}
|
||||
|
|
@ -162,7 +163,7 @@ func TestImportCommand_RunValue(t *testing.T) {
|
|||
buf := bytes.Buffer{}
|
||||
stdin, stdout, stderr := GetIO(buf)
|
||||
cm := NewImportCommand(stdin, stdout, stderr)
|
||||
file, err := ioutil.TempFile("", "import-value.csv")
|
||||
file, err := testhook.TempFile(t, "import-value.csv")
|
||||
if err != nil {
|
||||
t.Fatalf("creating tempfile: %v", err)
|
||||
}
|
||||
|
|
@ -207,7 +208,7 @@ func TestImportCommand_RunKeys(t *testing.T) {
|
|||
buf := bytes.Buffer{}
|
||||
stdin, stdout, stderr := GetIO(buf)
|
||||
cm := NewImportCommand(stdin, stdout, stderr)
|
||||
file, err := ioutil.TempFile("", "import-key.csv")
|
||||
file, err := testhook.TempFile(t, "import-key.csv")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -247,7 +248,7 @@ func TestImportCommand_KeyReplication(t *testing.T) {
|
|||
buf := bytes.Buffer{}
|
||||
stdin, stdout, stderr := GetIO(buf)
|
||||
cm := NewImportCommand(stdin, stdout, stderr)
|
||||
file, err := ioutil.TempFile("", "import-key.csv")
|
||||
file, err := testhook.TempFile(t, "import-key.csv")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -327,7 +328,7 @@ func TestImportCommand_RunValueKeys(t *testing.T) {
|
|||
buf := bytes.Buffer{}
|
||||
stdin, stdout, stderr := GetIO(buf)
|
||||
cm := NewImportCommand(stdin, stdout, stderr)
|
||||
file, err := ioutil.TempFile("", "import-key.csv")
|
||||
file, err := testhook.TempFile(t, "import-key.csv")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -373,7 +374,7 @@ func TestImportCommand_InvalidFile(t *testing.T) {
|
|||
cm.Host = cmd.API.Node().URI.HostPort()
|
||||
cm.Index = "i"
|
||||
cm.Field = "f"
|
||||
file, err := ioutil.TempFile("", "import.csv")
|
||||
file, err := testhook.TempFile(t, "import.csv")
|
||||
if err != nil {
|
||||
t.Fatalf("creating tempfile: %v", err)
|
||||
}
|
||||
|
|
@ -387,7 +388,7 @@ func TestImportCommand_InvalidFile(t *testing.T) {
|
|||
t.Fatalf("expect error: invalid row id on row, actual: %s", err)
|
||||
}
|
||||
|
||||
file, err = ioutil.TempFile("", "import1.csv")
|
||||
file, err = testhook.TempFile(t, "import1.csv")
|
||||
if err != nil {
|
||||
t.Fatalf("creating tempfile: %v", err)
|
||||
}
|
||||
|
|
@ -401,7 +402,7 @@ func TestImportCommand_InvalidFile(t *testing.T) {
|
|||
t.Fatalf("expect error: invalid column id on row, actual: %s", err)
|
||||
}
|
||||
|
||||
file, err = ioutil.TempFile("", "import1.csv")
|
||||
file, err = testhook.TempFile(t, "import1.csv")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -415,7 +416,7 @@ func TestImportCommand_InvalidFile(t *testing.T) {
|
|||
t.Fatalf("expect error: invalid timestamp on row, actual: %s", err)
|
||||
}
|
||||
|
||||
file, err = ioutil.TempFile("", "import1.csv")
|
||||
file, err = testhook.TempFile(t, "import1.csv")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -458,7 +459,7 @@ func TestImportCommand_BugOverwriteValue(t *testing.T) {
|
|||
buf := bytes.Buffer{}
|
||||
stdin, stdout, stderr := GetIO(buf)
|
||||
cm := NewImportCommand(stdin, stdout, stderr)
|
||||
file, err := ioutil.TempFile("", "import-value.csv")
|
||||
file, err := testhook.TempFile(t, "import-value.csv")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -490,7 +491,7 @@ func TestImportCommand_BugOverwriteValue(t *testing.T) {
|
|||
}
|
||||
|
||||
file.Close()
|
||||
file, err = ioutil.TempFile("", "import-value2.csv")
|
||||
file, err = testhook.TempFile(t, "import-value2.csv")
|
||||
if err != nil {
|
||||
t.Fatalf("Error creating tempfile: %s", err)
|
||||
}
|
||||
|
|
@ -505,7 +506,7 @@ func TestImportCommand_BugOverwriteValue(t *testing.T) {
|
|||
}
|
||||
|
||||
file.Close()
|
||||
file, err = ioutil.TempFile("", "import-value3.csv")
|
||||
file, err = testhook.TempFile(t, "import-value3.csv")
|
||||
if err != nil {
|
||||
t.Fatalf("Error creating tempfile: %s", err)
|
||||
}
|
||||
|
|
@ -547,7 +548,7 @@ func TestImportCommand_RunBool(t *testing.T) {
|
|||
cm.Field = "f"
|
||||
|
||||
t.Run("Valid", func(t *testing.T) {
|
||||
file, err := ioutil.TempFile("", "import-bool.csv")
|
||||
file, err := testhook.TempFile(t, "import-bool.csv")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -565,7 +566,7 @@ func TestImportCommand_RunBool(t *testing.T) {
|
|||
|
||||
// Ensure that invalid bool values return an error.
|
||||
t.Run("Invalid", func(t *testing.T) {
|
||||
file, err := ioutil.TempFile("", "import-invalid-bool.csv")
|
||||
file, err := testhook.TempFile(t, "import-invalid-bool.csv")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -18,10 +18,11 @@ import (
|
|||
"bytes"
|
||||
"context"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
)
|
||||
|
||||
func TestInspectCommand_Run(t *testing.T) {
|
||||
|
|
@ -30,7 +31,7 @@ func TestInspectCommand_Run(t *testing.T) {
|
|||
r, w, _ := os.Pipe()
|
||||
|
||||
cm := NewInspectCommand(stdin, w, w)
|
||||
file, err := ioutil.TempFile("", "inspectTest")
|
||||
file, err := testhook.TempFile(t, "inspectTest")
|
||||
if err != nil {
|
||||
t.Fatalf("Error creating tempfile: %s", err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -16,7 +16,6 @@ package pilosa
|
|||
|
||||
import (
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
|
|
@ -25,6 +24,7 @@ import (
|
|||
"github.com/pilosa/pilosa/v2/rbf"
|
||||
"github.com/pilosa/pilosa/v2/shardwidth"
|
||||
txkey "github.com/pilosa/pilosa/v2/short_txkey"
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
)
|
||||
|
||||
// Shard per db evaluation
|
||||
|
|
@ -70,7 +70,7 @@ func TestShardPerDB_SetBit(t *testing.T) {
|
|||
|
||||
// test that we find all *local* shards
|
||||
func Test_DBPerShard_GetShardsForIndex_LocalOnly(t *testing.T) {
|
||||
tmpdir, err := ioutil.TempDir("", "Test_DBPerShard_GetShardsForIndex_LocalOnly")
|
||||
tmpdir, err := testhook.TempDir(t, "Test_DBPerShard_GetShardsForIndex_LocalOnly")
|
||||
panicOn(err)
|
||||
defer os.RemoveAll(tmpdir)
|
||||
|
||||
|
|
@ -324,7 +324,7 @@ func makeTxTestDBWithViewsShards(holder *Holder, idx *Index, exp *FieldView2Shar
|
|||
|
||||
// test that rbf can give us a map[view]*shardSet
|
||||
func Test_DBPerShard_GetFieldView2Shards_map_from_RBF(t *testing.T) {
|
||||
tmpdir, err := ioutil.TempDir("", "Test_DBPerShard_GetFieldView2Shards_map_from_RBF")
|
||||
tmpdir, err := testhook.TempDir(t, "Test_DBPerShard_GetFieldView2Shards_map_from_RBF")
|
||||
panicOn(err)
|
||||
defer os.RemoveAll(tmpdir)
|
||||
|
||||
|
|
|
|||
|
|
@ -100,7 +100,7 @@ func (l *leasedKV) consumeLease(ch <-chan *clientv3.LeaseKeepAliveResponse) {
|
|||
return
|
||||
}
|
||||
|
||||
if ok := retry(1*time.Second, func() error {
|
||||
if e := retry("consumeLease", 1*time.Second, func() error {
|
||||
kaChann, err := l.create(l.value)
|
||||
if err != nil {
|
||||
return err
|
||||
|
|
@ -108,13 +108,13 @@ func (l *leasedKV) consumeLease(ch <-chan *clientv3.LeaseKeepAliveResponse) {
|
|||
|
||||
go l.consumeLease(kaChann)
|
||||
return nil
|
||||
}); !ok {
|
||||
log.Println("lease cannot be recreated. Key:", l.key)
|
||||
}); e != nil {
|
||||
log.Printf("lease %q cannot be recreated: %v", l.key, e)
|
||||
l.mu.Unlock()
|
||||
return
|
||||
}
|
||||
|
||||
log.Println("lease recreated after a problem. Key:", l.key)
|
||||
log.Printf("lease %q recreated after a problem", l.key)
|
||||
l.mu.Unlock()
|
||||
return
|
||||
}
|
||||
|
|
@ -172,20 +172,21 @@ func (l *leasedKV) Get(ctx context.Context) (string, error) {
|
|||
return l.value, nil
|
||||
}
|
||||
|
||||
func retry(sleep time.Duration, f func() error) bool {
|
||||
func retry(desc string, sleep time.Duration, f func() error) (err error) {
|
||||
for {
|
||||
err := f()
|
||||
if err == nil {
|
||||
return true
|
||||
lastErr := f()
|
||||
if lastErr == nil {
|
||||
return lastErr
|
||||
}
|
||||
|
||||
// sometimes the element in charge of stopping the lease renewal doesn't do it, causing context errors.
|
||||
if errors.Is(err, context.DeadlineExceeded) {
|
||||
return false
|
||||
if errors.Is(lastErr, context.DeadlineExceeded) {
|
||||
if err != nil {
|
||||
return err
|
||||
} else {
|
||||
return lastErr
|
||||
}
|
||||
}
|
||||
|
||||
log.Printf("%s: got error %v, retrying", desc, lastErr)
|
||||
err = lastErr
|
||||
time.Sleep(sleep)
|
||||
|
||||
log.Println("retrying after error:", err)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -17,12 +17,12 @@ package etcd
|
|||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/disco"
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
"go.etcd.io/etcd/embed"
|
||||
"go.etcd.io/etcd/etcdserver/api/v3client"
|
||||
)
|
||||
|
|
@ -33,7 +33,7 @@ const newVal = "newValue"
|
|||
func TestLeasedKv(t *testing.T) {
|
||||
cfg := embed.NewConfig()
|
||||
|
||||
dir, err := ioutil.TempDir("", "leasedkv-*")
|
||||
dir, err := testhook.TempDir(t, "leasedkv-*")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
|
|||
135
executor.go
135
executor.go
|
|
@ -15,6 +15,7 @@
|
|||
package pilosa
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
|
|
@ -3162,6 +3163,16 @@ type FieldRow struct {
|
|||
Value *int64 `json:"value,omitempty"`
|
||||
}
|
||||
|
||||
func (p FieldRow) Less(o FieldRow) bool {
|
||||
if p.Field == o.Field {
|
||||
if p.RowKey != "" {
|
||||
return p.RowKey < o.RowKey
|
||||
}
|
||||
return p.RowID < o.RowID
|
||||
}
|
||||
return p.Field < o.Field
|
||||
}
|
||||
|
||||
func (fr *FieldRow) Clone() (clone *FieldRow) {
|
||||
clone = &FieldRow{
|
||||
Field: fr.Field,
|
||||
|
|
@ -3328,7 +3339,7 @@ func (g *GroupCounts) ToRows(callback func(*pb.RowResponse) error) error {
|
|||
// customizes the JSON output of the aggregate field label.
|
||||
func (g *GroupCounts) MarshalJSON() ([]byte, error) {
|
||||
groups := g.Groups()
|
||||
var counts interface{} = groups
|
||||
var counts interface{} = SortedGroupCount(groups)
|
||||
|
||||
if len(groups) == 0 {
|
||||
return []byte("[]"), nil
|
||||
|
|
@ -3342,6 +3353,50 @@ func (g *GroupCounts) MarshalJSON() ([]byte, error) {
|
|||
return json.Marshal(counts)
|
||||
}
|
||||
|
||||
type SortedGroupCount []GroupCount
|
||||
|
||||
// Provides determanistic ordering for the JSON output
|
||||
func (sortgroup SortedGroupCount) MarshalJSON() ([]byte, error) {
|
||||
sort.Slice(sortgroup, func(i, j int) bool {
|
||||
if sortgroup[j].Count == sortgroup[i].Count {
|
||||
for x := range sortgroup[j].Group {
|
||||
if sortgroup[j].Group[x] == sortgroup[i].Group[x] {
|
||||
continue
|
||||
}
|
||||
return sortgroup[j].Group[x].Less(sortgroup[i].Group[x])
|
||||
}
|
||||
|
||||
}
|
||||
return sortgroup[j].Count < sortgroup[i].Count
|
||||
})
|
||||
buffer := new(bytes.Buffer)
|
||||
_, err := buffer.Write([]byte("["))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for i, groupCount := range sortgroup {
|
||||
jsonValue, err := json.Marshal(groupCount)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if i != 0 {
|
||||
_, err = buffer.Write([]byte(","))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
_, err = buffer.Write(jsonValue)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
_, err = buffer.Write([]byte("]"))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return buffer.Bytes(), nil
|
||||
}
|
||||
|
||||
// GroupCount represents a result item for a group by query.
|
||||
type GroupCount struct {
|
||||
Group []FieldRow `json:"group"`
|
||||
|
|
@ -3349,6 +3404,32 @@ type GroupCount struct {
|
|||
Agg int64 `json:"-"`
|
||||
}
|
||||
|
||||
// comparing json structures in the oracle demand that map values have a consistent order
|
||||
// docs say it should be lexigraphically but experimentation proved otherwise int the
|
||||
// groupcount state
|
||||
func (group GroupCount) MarshalJSON() ([]byte, error) {
|
||||
// TODO(twg) figure out how to access struct tags from code
|
||||
var buf bytes.Buffer
|
||||
sort.Slice(group.Group, func(i, j int) bool {
|
||||
return group.Group[i].Field < group.Group[j].Field
|
||||
})
|
||||
|
||||
//write the Counts
|
||||
buf.WriteString(fmt.Sprintf(`{"count": %v,`, group.Count))
|
||||
//write the Groups
|
||||
buf.WriteString(`"group":[`)
|
||||
for i, g := range group.Group {
|
||||
part, _ := g.MarshalJSON()
|
||||
if i != 0 {
|
||||
buf.Write([]byte(","))
|
||||
}
|
||||
buf.Write(part)
|
||||
}
|
||||
|
||||
buf.WriteString("]}")
|
||||
return buf.Bytes(), nil
|
||||
}
|
||||
|
||||
type groupCountSum struct {
|
||||
Group []FieldRow `json:"group"`
|
||||
Count uint64 `json:"count"`
|
||||
|
|
@ -3845,6 +3926,43 @@ type ExtractedTableColumn struct {
|
|||
Rows []interface{} `json:"rows"`
|
||||
}
|
||||
|
||||
func (etc ExtractedTableColumn) MarshalJSON() ([]byte, error) {
|
||||
buf := new(bytes.Buffer)
|
||||
buf.WriteString(`{"column":`)
|
||||
b, _ := etc.Column.MarshalJSON()
|
||||
buf.Write(b)
|
||||
buf.WriteString(`,"rows":[`)
|
||||
for i, row := range etc.Rows {
|
||||
if i != 0 {
|
||||
buf.WriteByte(44) //write a comma
|
||||
}
|
||||
buf.WriteString(`[`)
|
||||
switch v := row.(type) {
|
||||
case []string:
|
||||
sort.Slice(v, func(i, j int) bool {
|
||||
return v[i] < v[j]
|
||||
})
|
||||
for si, s := range v {
|
||||
if si != 0 {
|
||||
buf.WriteByte(44) //write a comma
|
||||
}
|
||||
buf.WriteString(fmt.Sprintf(`"%v"`, s))
|
||||
}
|
||||
case []uint64:
|
||||
for si, s := range v {
|
||||
if si != 0 {
|
||||
buf.WriteByte(44) //write a comma
|
||||
}
|
||||
buf.WriteString(fmt.Sprintf(`%v`, s))
|
||||
}
|
||||
}
|
||||
buf.WriteString(`]`)
|
||||
|
||||
}
|
||||
buf.WriteString(`]}`)
|
||||
return buf.Bytes(), nil
|
||||
}
|
||||
|
||||
type ExtractedTable struct {
|
||||
Fields []ExtractedTableField `json:"fields"`
|
||||
Columns []ExtractedTableColumn `json:"columns"`
|
||||
|
|
@ -7407,6 +7525,21 @@ type ValCount struct {
|
|||
Count int64 `json:"count"`
|
||||
}
|
||||
|
||||
func (vc ValCount) MarshalJSON() ([]byte, error) {
|
||||
buf := new(bytes.Buffer)
|
||||
nullOrVal := "null"
|
||||
if vc.DecimalVal != nil {
|
||||
nullOrVal = fmt.Sprintf("%v", vc.DecimalVal)
|
||||
}
|
||||
_, err := buf.WriteString(fmt.Sprintf(`{
|
||||
"count": %v,
|
||||
"decimalValue": %v,
|
||||
"floatValue": %v,
|
||||
"value":%v
|
||||
}`, vc.Count, nullOrVal, vc.FloatVal, vc.Val))
|
||||
return buf.Bytes(), err
|
||||
}
|
||||
|
||||
func (v *ValCount) Clone() (r *ValCount) {
|
||||
r = &ValCount{
|
||||
Val: v.Val,
|
||||
|
|
|
|||
|
|
@ -3166,7 +3166,7 @@ func BenchmarkImportIntoLargeFragment(b *testing.B) {
|
|||
if err != nil {
|
||||
b.Fatalf("opening frag file: %v", err)
|
||||
}
|
||||
fi, err := ioutil.TempFile(*TempDir, "")
|
||||
fi, err := testhook.TempFileInDir(b, *TempDir, "")
|
||||
if err != nil {
|
||||
b.Fatalf("getting temp file: %v", err)
|
||||
}
|
||||
|
|
@ -3216,7 +3216,7 @@ func BenchmarkImportRoaringIntoLargeFragment(b *testing.B) {
|
|||
if err != nil {
|
||||
b.Fatalf("opening frag file: %v", err)
|
||||
}
|
||||
fi, err := ioutil.TempFile(*TempDir, "")
|
||||
fi, err := testhook.TempFileInDir(b, *TempDir, "")
|
||||
if err != nil {
|
||||
b.Fatalf("getting temp file: %v", err)
|
||||
}
|
||||
|
|
@ -3472,6 +3472,10 @@ 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++ {
|
||||
// DO NOT CONVERT THIS ONE TO USE TESTHOOK.
|
||||
// We're deleting these files as we go because
|
||||
// otherwise the benchmark could fill up
|
||||
// $TMPDIR before it finishes running.
|
||||
f, err := ioutil.TempFile(*TempDir, "")
|
||||
if err != nil {
|
||||
b.Fatalf("getting temp file: %v", err)
|
||||
|
|
@ -3479,14 +3483,17 @@ func BenchmarkFileWrite(b *testing.B) {
|
|||
b.StartTimer()
|
||||
_, err = f.Write(data)
|
||||
if err != nil {
|
||||
os.Remove(f.Name())
|
||||
b.Fatal(err)
|
||||
}
|
||||
err = f.Sync()
|
||||
if err != nil {
|
||||
os.Remove(f.Name())
|
||||
b.Fatal(err)
|
||||
}
|
||||
err = f.Close()
|
||||
if err != nil {
|
||||
os.Remove(f.Name())
|
||||
b.Fatal(err)
|
||||
}
|
||||
b.StopTimer()
|
||||
|
|
|
|||
7
go.mod
7
go.mod
|
|
@ -16,7 +16,7 @@ require (
|
|||
github.com/glycerine/goconvey v0.0.0-20190410193231-58a59202ab31 // indirect
|
||||
github.com/glycerine/idem v0.0.0-20190127113923-7a8083893311
|
||||
github.com/gogo/protobuf v1.2.1
|
||||
github.com/golang/protobuf v1.4.2
|
||||
github.com/golang/protobuf v1.3.3
|
||||
github.com/google/go-cmp v0.5.2
|
||||
github.com/google/uuid v1.1.4 // indirect
|
||||
github.com/gopherjs/gopherjs v0.0.0-20200217142428-fce0ec30dd00 // indirect
|
||||
|
|
@ -34,7 +34,6 @@ require (
|
|||
github.com/prometheus/client_model v0.1.0
|
||||
github.com/prometheus/prom2json v1.3.0
|
||||
github.com/rakyll/statik v0.1.7
|
||||
github.com/remyoudompheng/bigfft v0.0.0-20190728182440-6a916e37a237 // indirect
|
||||
github.com/rs/cors v1.7.0 // indirect
|
||||
github.com/satori/go.uuid v1.2.0
|
||||
github.com/shirou/gopsutil/v3 v3.20.11
|
||||
|
|
@ -52,12 +51,12 @@ require (
|
|||
golang.org/x/mod v0.3.1-0.20200828183125-ce943fd02449
|
||||
golang.org/x/net v0.0.0-20200822124328-c89045814202 // indirect
|
||||
golang.org/x/sync v0.0.0-20200625203802-6e8e738ad208
|
||||
golang.org/x/sys v0.0.0-20201214095126-aec9a390925b // indirect
|
||||
golang.org/x/sys v0.0.0-20210124154548-22da62e12c0c // indirect
|
||||
golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1 // indirect
|
||||
google.golang.org/grpc v1.28.0
|
||||
gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f // indirect
|
||||
gopkg.in/yaml.v2 v2.3.0 // indirect
|
||||
modernc.org/mathutil v1.0.0
|
||||
modernc.org/mathutil v1.2.2
|
||||
modernc.org/strutil v1.0.0
|
||||
sigs.k8s.io/yaml v1.2.0 // indirect
|
||||
vitess.io/vitess v3.0.0-rc.3.0.20190602171040-12bfde34629c+incompatible
|
||||
|
|
|
|||
57
go.sum
57
go.sum
|
|
@ -28,7 +28,6 @@ github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf/go.mod h1:ybxpYRF
|
|||
github.com/alecthomas/units v0.0.0-20190717042225-c3de453c63f4/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0=
|
||||
github.com/alecthomas/units v0.0.0-20190924025748-f65c72e2690d/go.mod h1:rBZYJk541a8SKzHPHnH3zbiI+7dagKZ0cgpgrD7Fyho=
|
||||
github.com/armon/circbuf v0.0.0-20150827004946-bbbad097214e/go.mod h1:3U/XgcO3hCbHZ8TKRvWD2dDTCfh9M9ya+I9JpbB7O8o=
|
||||
github.com/armon/go-metrics v0.0.0-20180917152333-f0300d1749da h1:8GUt8eRujhVEGZFFEjBj46YV4rDjvGrNxb0KMWYkL2I=
|
||||
github.com/armon/go-metrics v0.0.0-20180917152333-f0300d1749da/go.mod h1:Q73ZrmVTwzkszR9V5SSuryQ31EELlFMUz1kKyl939pY=
|
||||
github.com/armon/go-radix v0.0.0-20180808171621-7fddfc383310/go.mod h1:ufUuZ+zHj4x4TnLV4JWEpy2hxWSpsRywHrMgIH9cCH8=
|
||||
github.com/beevik/ntp v0.3.0 h1:xzVrPrE4ziasFXgBVBZJDP0Wg/KpMwk2KHJ4Ba8GrDw=
|
||||
|
|
@ -49,9 +48,7 @@ github.com/cockroachdb/datadriven v0.0.0-20190809214429-80d97fb3cbaa h1:OaNxuTZr
|
|||
github.com/cockroachdb/datadriven v0.0.0-20190809214429-80d97fb3cbaa/go.mod h1:zn76sxSg3SzpJ0PPJaLDCu+Bu0Lg3sKTORVIj19EIF8=
|
||||
github.com/codahale/hdrhistogram v0.0.0-20161010025455-3a0bb77429bd h1:qMd81Ts1T2OTKmB4acZcyKaMtRnY5Y44NuXGX2GFJ1w=
|
||||
github.com/codahale/hdrhistogram v0.0.0-20161010025455-3a0bb77429bd/go.mod h1:sE/e/2PUdi/liOCUjSTXgM1o87ZssimdTWN964YiIeI=
|
||||
github.com/coreos/bbolt v1.3.2 h1:wZwiHHUieZCquLkDL0B8UhzreNWsPHooDAG3q34zk0s=
|
||||
github.com/coreos/bbolt v1.3.2/go.mod h1:iRUV2dpdMOn7Bo10OQBFzIJO9kkE559Wcmn+qkEiiKk=
|
||||
github.com/coreos/etcd v3.3.13+incompatible h1:8F3hqu9fGYLBifCmRCJsicFqDx/D68Rt3q1JMazcgBQ=
|
||||
github.com/coreos/etcd v3.3.13+incompatible/go.mod h1:uF7uidLiAD3TWHmW31ZFd/JWoc32PjwdhPthX9715RE=
|
||||
github.com/coreos/go-semver v0.2.0/go.mod h1:nnelYz7RCh+5ahJtPPxZlU+153eP4D4r3EedlOD2RNk=
|
||||
github.com/coreos/go-semver v0.3.0 h1:wkHLiw0WNATZnSG7epLsujiMCgPAc9xhjJ4tgnAxmfM=
|
||||
|
|
@ -81,7 +78,6 @@ github.com/envoyproxy/go-control-plane v0.9.1-0.20191026205805-5f8ba28d4473/go.m
|
|||
github.com/envoyproxy/go-control-plane v0.9.4/go.mod h1:6rpuAdCZL397s3pYoYcLgu1mIlRU8Am5FuJP05cCM98=
|
||||
github.com/envoyproxy/protoc-gen-validate v0.1.0/go.mod h1:iSmxcyjqTsJpI2R4NaDN7+kN2VEUnK/pcBlmesArF7c=
|
||||
github.com/fatih/color v1.7.0/go.mod h1:Zm6kSWBoL9eyXnKyktHP6abPY2pDugNf5KwzbycvMj4=
|
||||
github.com/fsnotify/fsnotify v1.4.7 h1:IXs+QLmnXW2CcXuY+8Mzv/fWEsPGWxqefPtCP5CnV9I=
|
||||
github.com/fsnotify/fsnotify v1.4.7/go.mod h1:jwhsz4b93w/PPRr/qN1Yymfu8t87LnFCMoQvtojpjFo=
|
||||
github.com/fsnotify/fsnotify v1.4.9 h1:hsms1Qyu0jgnwNXIxa+/V/PDsU6CfLf6CNO8H7IWoS4=
|
||||
github.com/fsnotify/fsnotify v1.4.9/go.mod h1:znqG4EE+3YCdAaPaxE2ZRY/06pZUdp0tY4IgpuI1SZQ=
|
||||
|
|
@ -115,21 +111,11 @@ github.com/golang/protobuf v1.3.1/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5y
|
|||
github.com/golang/protobuf v1.3.2/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U=
|
||||
github.com/golang/protobuf v1.3.3 h1:gyjaxf+svBWX08ZjK86iN9geUJF0H6gp2IRKX6Nf6/I=
|
||||
github.com/golang/protobuf v1.3.3/go.mod h1:vzj43D7+SQXF/4pzW/hwtAqwc6iTitCiVSaWz5lYuqw=
|
||||
github.com/golang/protobuf v1.4.0-rc.1/go.mod h1:ceaxUfeHdC40wWswd/P6IGgMaK3YpKi5j83Wpe3EHw8=
|
||||
github.com/golang/protobuf v1.4.0-rc.1.0.20200221234624-67d41d38c208/go.mod h1:xKAWHe0F5eneWXFV3EuXVDTCmh+JuBKY0li0aMyXATA=
|
||||
github.com/golang/protobuf v1.4.0-rc.2/go.mod h1:LlEzMj4AhA7rCAGe4KMBDvJI+AwstrUpVNzEA03Pprs=
|
||||
github.com/golang/protobuf v1.4.0-rc.4.0.20200313231945-b860323f09d0/go.mod h1:WU3c8KckQ9AFe+yFwt9sWVRKCVIyN9cPHBJSNnbL67w=
|
||||
github.com/golang/protobuf v1.4.0/go.mod h1:jodUvKwWbYaEsadDk5Fwe5c77LiNKVO9IDvqG2KuDX0=
|
||||
github.com/golang/protobuf v1.4.2 h1:+Z5KGCizgyZCbGh1KZqA0fcLLkwbsjIzS4aV2v7wJX0=
|
||||
github.com/golang/protobuf v1.4.2/go.mod h1:oDoupMAO8OvCJWAcko0GGGIgR6R6ocIYbsSw735rRwI=
|
||||
github.com/google/btree v0.0.0-20180813153112-4030bb1f1f0c/go.mod h1:lNA+9X1NB3Zf8V7Ke586lFgjr2dZNuvo3lPJSGZ5JPQ=
|
||||
github.com/google/btree v1.0.0 h1:0udJVsspx3VBr5FwtLhQQtuAsVc79tTq0ocGIPAU6qo=
|
||||
github.com/google/btree v1.0.0/go.mod h1:lNA+9X1NB3Zf8V7Ke586lFgjr2dZNuvo3lPJSGZ5JPQ=
|
||||
github.com/google/go-cmp v0.2.0/go.mod h1:oXzfMopK8JAjlY9xF4vHSVASa0yLyX7SntLO5aqRK0M=
|
||||
github.com/google/go-cmp v0.3.0/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU=
|
||||
github.com/google/go-cmp v0.3.1/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU=
|
||||
github.com/google/go-cmp v0.4.0 h1:xsAVV57WRhGj6kEIi8ReJzQlHHqcBYCElAvkovg3B/4=
|
||||
github.com/google/go-cmp v0.4.0/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE=
|
||||
github.com/google/go-cmp v0.5.2 h1:X2ev0eStA3AbceY54o37/0PQ/UWqKEiiO2dKL5OPaFM=
|
||||
github.com/google/go-cmp v0.5.2/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE=
|
||||
github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
|
||||
|
|
@ -152,39 +138,28 @@ github.com/gorilla/mux v1.7.0/go.mod h1:1lud6UwP+6orDFRuTfBEV8e9/aOM/c4fVVCaMa2z
|
|||
github.com/gorilla/websocket v0.0.0-20170926233335-4201258b820c/go.mod h1:E7qHFY5m1UJ88s3WnNqhKjPHQ0heANvMoAMk2YaljkQ=
|
||||
github.com/gorilla/websocket v1.4.2 h1:+/TMaTYc4QFitKJxsQ7Yye35DkWvkdLcvGKqM+x0Ufc=
|
||||
github.com/gorilla/websocket v1.4.2/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE=
|
||||
github.com/grpc-ecosystem/go-grpc-middleware v1.0.0 h1:Iju5GlWwrvL6UBg4zJJt3btmonfrMlCDdsejg4CZE7c=
|
||||
github.com/grpc-ecosystem/go-grpc-middleware v1.0.0/go.mod h1:FiyG127CGDf3tlThmgyCl78X/SZQqEOJBCDaAfeWzPs=
|
||||
github.com/grpc-ecosystem/go-grpc-middleware v1.0.1-0.20190118093823-f849b5445de4 h1:z53tR0945TRRQO/fLEVPI6SMv7ZflF0TEaTAoU7tOzg=
|
||||
github.com/grpc-ecosystem/go-grpc-middleware v1.0.1-0.20190118093823-f849b5445de4/go.mod h1:FiyG127CGDf3tlThmgyCl78X/SZQqEOJBCDaAfeWzPs=
|
||||
github.com/grpc-ecosystem/go-grpc-prometheus v1.2.0 h1:Ovs26xHkKqVztRpIrF/92BcuyuQ/YW4NSIpoGtfXNho=
|
||||
github.com/grpc-ecosystem/go-grpc-prometheus v1.2.0/go.mod h1:8NvIoxWQoOIhqOTXgfV/d3M/q6VIi02HzZEHgUlZvzk=
|
||||
github.com/grpc-ecosystem/grpc-gateway v1.9.0 h1:bM6ZAFZmc/wPFaRDi0d5L7hGEZEx/2u+Tmr2evNHDiI=
|
||||
github.com/grpc-ecosystem/grpc-gateway v1.9.0/go.mod h1:vNeuVxBJEsws4ogUvrchl83t/GYV9WGTSLVdBhOQFDY=
|
||||
github.com/grpc-ecosystem/grpc-gateway v1.9.5 h1:UImYN5qQ8tuGpGE16ZmjvcTtTw24zw1QAp/SlnNrZhI=
|
||||
github.com/grpc-ecosystem/grpc-gateway v1.9.5/go.mod h1:vNeuVxBJEsws4ogUvrchl83t/GYV9WGTSLVdBhOQFDY=
|
||||
github.com/hashicorp/consul/api v1.1.0/go.mod h1:VmuI/Lkw1nC05EYQWNKwWGbkg+FbDBtguAZLlVdkD9Q=
|
||||
github.com/hashicorp/consul/sdk v0.1.1/go.mod h1:VKf9jXwCTEY1QZP2MOLRhb5i/I/ssyNV1vwHyQBF0x8=
|
||||
github.com/hashicorp/errwrap v1.0.0 h1:hLrqtEDnRye3+sgx6z4qVLNuviH3MR5aQ0ykNJa/UYA=
|
||||
github.com/hashicorp/errwrap v1.0.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4=
|
||||
github.com/hashicorp/go-cleanhttp v0.5.1/go.mod h1:JpRdi6/HCYpAwUzNwuwqhbovhLtngrth3wmdIIUrZ80=
|
||||
github.com/hashicorp/go-immutable-radix v1.0.0 h1:AKDB1HM5PWEA7i4nhcpwOrO2byshxBjXVn/J/3+z5/0=
|
||||
github.com/hashicorp/go-immutable-radix v1.0.0/go.mod h1:0y9vanUI8NX6FsYoO3zeMjhV/C5i9g4Q3DwcSNZ4P60=
|
||||
github.com/hashicorp/go-msgpack v0.5.3 h1:zKjpN5BK/P5lMYrLmBHdBULWbJ0XpYR+7NGzqkZzoD4=
|
||||
github.com/hashicorp/go-msgpack v0.5.3/go.mod h1:ahLV/dePpqEmjfWmKiqvPkv/twdG7iPBM1vqhUKIvfM=
|
||||
github.com/hashicorp/go-multierror v1.0.0 h1:iVjPR7a6H0tWELX5NxNe7bYopibicUzc7uPribsnS6o=
|
||||
github.com/hashicorp/go-multierror v1.0.0/go.mod h1:dHtQlpGsu+cZNNAkkCN/P3hoUDHhCYQXV3UM06sGGrk=
|
||||
github.com/hashicorp/go-rootcerts v1.0.0/go.mod h1:K6zTfqpRlCUIjkwsN4Z+hiSfzSTQa6eBIzfwKfwNnHU=
|
||||
github.com/hashicorp/go-sockaddr v1.0.0 h1:GeH6tui99pF4NJgfnhp+L6+FfobzVW3Ah46sLo0ICXs=
|
||||
github.com/hashicorp/go-sockaddr v1.0.0/go.mod h1:7Xibr9yA9JjQq1JpNB2Vw7kxv8xerXegt+ozgdvDeDU=
|
||||
github.com/hashicorp/go-syslog v1.0.0/go.mod h1:qPfqrKkXGihmCqbJM2mZgkZGvKG1dFdvsLplgctolz4=
|
||||
github.com/hashicorp/go-uuid v1.0.0 h1:RS8zrF7PhGwyNPOtxSClXXj9HA8feRnJzgnI1RJCSnM=
|
||||
github.com/hashicorp/go-uuid v1.0.0/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro=
|
||||
github.com/hashicorp/go-uuid v1.0.1 h1:fv1ep09latC32wFoVwnqcnKJGnMSdBanPczbHAYm1BE=
|
||||
github.com/hashicorp/go-uuid v1.0.1/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro=
|
||||
github.com/hashicorp/go.net v0.0.1/go.mod h1:hjKkEWcCURg++eb33jQU7oqQcI9XDCnUzHA0oac0k90=
|
||||
github.com/hashicorp/golang-lru v0.5.0 h1:CL2msUPvZTLb5O648aiLNJw3hnBxN2+1Jq8rCOH9wdo=
|
||||
github.com/hashicorp/golang-lru v0.5.0/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8=
|
||||
github.com/hashicorp/golang-lru v0.5.1 h1:0hERBMJE1eitiLkihrMvRVBYAkpHzc/J3QdDN+dAcgU=
|
||||
github.com/hashicorp/golang-lru v0.5.1/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8=
|
||||
github.com/hashicorp/hcl v1.0.0 h1:0Anlzjpi4vEasTeNFn2mLJgTSwt0+6sfsiTG8qcWGx4=
|
||||
github.com/hashicorp/hcl v1.0.0/go.mod h1:E5yfLk+7swimpb2L/Alb/PJmXilQ/rhwaUYs4T20WEQ=
|
||||
|
|
@ -198,7 +173,6 @@ github.com/inconshreveable/mousetrap v1.0.0 h1:Z8tu5sraLXCXIcARxBp/8cbvlwVa7Z1NH
|
|||
github.com/inconshreveable/mousetrap v1.0.0/go.mod h1:PxqpIevigyE2G7u3NXJIT2ANytuPF1OarO4DADm73n8=
|
||||
github.com/jonboulle/clockwork v0.1.0 h1:VKV+ZcuP6l3yW9doeqz6ziZGgcynBVQO+obU0+0hcPo=
|
||||
github.com/jonboulle/clockwork v0.1.0/go.mod h1:Ii8DK3G1RaLaWxj9trq07+26W01tbo22gdxWY5EU2bo=
|
||||
github.com/json-iterator/go v1.1.6 h1:MrUvLMLTMxbqFJ9kzlvat/rYZqZnW3u4wkLzWTaFwKs=
|
||||
github.com/json-iterator/go v1.1.6/go.mod h1:+SdeFBvtyEkXs7REEP0seUULqWtbJapLOCVDaaPEHmU=
|
||||
github.com/json-iterator/go v1.1.7 h1:KfgG9LzI+pYjr4xvmz/5H4FXjokeP+rlHLhv3iH62Fo=
|
||||
github.com/json-iterator/go v1.1.7/go.mod h1:KdQUCv79m/52Kvf8AW2vK1V8akMuk1QjK/uOdHXbAo4=
|
||||
|
|
@ -207,16 +181,13 @@ github.com/jtolds/gls v4.20.0+incompatible h1:xdiiI2gbIgH/gLH7ADydsJ1uDOEzR8yvV7
|
|||
github.com/jtolds/gls v4.20.0+incompatible/go.mod h1:QJZ7F/aHp+rZTRtaJ1ow/lLfFfVYBRgL+9YlvaHOwJU=
|
||||
github.com/julienschmidt/httprouter v1.2.0/go.mod h1:SYymIcj16QtmaHHD7aYtjjsJG7VTCxuUUipMqKk8s4w=
|
||||
github.com/kisielk/errcheck v1.1.0/go.mod h1:EZBBE59ingxPouuu3KfxchcWSUPOHkagtvWXihfKN4Q=
|
||||
github.com/kisielk/gotool v1.0.0 h1:AV2c/EiW3KqPNT9ZKl07ehoAGi4C5/01Cfbblndcapg=
|
||||
github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck=
|
||||
github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ=
|
||||
github.com/konsorten/go-windows-terminal-sequences v1.0.2 h1:DB17ag19krx9CFsz4o3enTrPXyIXCl+2iCXH/aMAp9s=
|
||||
github.com/konsorten/go-windows-terminal-sequences v1.0.2/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ=
|
||||
github.com/kr/logfmt v0.0.0-20140226030751-b84e30acd515/go.mod h1:+0opPa2QZZtGFBFZlji/RkVcI2GknAs/DXo4wKdlNEc=
|
||||
github.com/kr/pretty v0.1.0 h1:L/CwN0zerZDmRFUapSPitk6f+Q3+0za1rQkzVuMiMFI=
|
||||
github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo=
|
||||
github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ=
|
||||
github.com/kr/text v0.1.0 h1:45sCR5RtlFHMR4UwH9sdQ5TC8v0qDQCHnXt+kaKSTVE=
|
||||
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
|
||||
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
|
||||
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
|
||||
|
|
@ -230,11 +201,9 @@ github.com/mattn/go-isatty v0.0.4/go.mod h1:M+lRXTBqGeGNdLjl/ufCoiOlB5xdOkqRJdNx
|
|||
github.com/mattn/go-runewidth v0.0.2/go.mod h1:LwmH8dsx7+W8Uxz3IHJYH5QSwggIsqBzpuz5H//U1FU=
|
||||
github.com/matttproud/golang_protobuf_extensions v1.0.1 h1:4hp9jkHxhMHkqkrB3Ix0jegS5sx/RkqARlsWZ6pIwiU=
|
||||
github.com/matttproud/golang_protobuf_extensions v1.0.1/go.mod h1:D8He9yQNgCq6Z5Ld7szi9bcBfOoFv/3dc6xSMkL2PC0=
|
||||
github.com/miekg/dns v1.0.14 h1:9jZdLNd/P4+SfEJ0TNyxYpsK8N4GtfylBLqtbYN1sbA=
|
||||
github.com/miekg/dns v1.0.14/go.mod h1:W1PPwlIAgtquWBMBEV9nkV9Cazfe8ScdGz/Lj7v3Nrg=
|
||||
github.com/mitchellh/cli v1.0.0/go.mod h1:hNIlj7HEI86fIcpObd7a0FcrxTWetlwJDGcceTlRvqc=
|
||||
github.com/mitchellh/go-homedir v1.0.0/go.mod h1:SfyaCUpYCn1Vlf4IUYiD9fPX4A5wJrkLzIz1N1q0pr0=
|
||||
github.com/mitchellh/go-homedir v1.1.0 h1:lukF9ziXFxDFPkA1vsr5zpc1XuPDn/wFntq5mG+4E0Y=
|
||||
github.com/mitchellh/go-homedir v1.1.0/go.mod h1:SfyaCUpYCn1Vlf4IUYiD9fPX4A5wJrkLzIz1N1q0pr0=
|
||||
github.com/mitchellh/go-testing-interface v1.0.0/go.mod h1:kRemZodwjscx+RGhAo8eIhFbs2+BFgRtFPeD/KE+zxI=
|
||||
github.com/mitchellh/gox v0.4.0/go.mod h1:Sd9lOJ0+aimLBi73mGofS1ycjY8lL3uZM3JPS42BGNg=
|
||||
|
|
@ -260,7 +229,6 @@ github.com/oklog/ulid v1.3.1/go.mod h1:CirwcVhetQ6Lv90oh/F+FBtV6XMibvdAFo93nm5qn
|
|||
github.com/olekukonko/tablewriter v0.0.0-20170122224234-a0225b3f23b5/go.mod h1:vsDQFd/mU46D+Z4whnwzcISnGGzXWMclvtLoiIKAKIo=
|
||||
github.com/opentracing/opentracing-go v1.1.0 h1:pWlfV3Bxv7k65HYwkikxat0+s3pV4bsqf19k25Ur8rU=
|
||||
github.com/opentracing/opentracing-go v1.1.0/go.mod h1:UkNAQd3GIcIGf0SeVgPpRdFStlNbqXla1AfSYxPUl2o=
|
||||
github.com/pascaldekloe/goe v0.0.0-20180627143212-57f6aae5913c h1:Lgl0gzECD8GnQ5QCWA8o6BtfL6mDH5rQgM4/fX3avOs=
|
||||
github.com/pascaldekloe/goe v0.0.0-20180627143212-57f6aae5913c/go.mod h1:lzWF7FIEvWOWxwDKqyGYQf6ZUaNfKdP144TG7ZOy1lc=
|
||||
github.com/pelletier/go-toml v1.2.0 h1:T5zMGML61Wp+FlcbWjRDT7yAxhJNAiPPLOFECq181zc=
|
||||
github.com/pelletier/go-toml v1.2.0/go.mod h1:5z9KED0ma1S8pY6P1sdut58dfprrGBbd/94hg7ilaic=
|
||||
|
|
@ -294,8 +262,8 @@ github.com/prometheus/prom2json v1.3.0/go.mod h1:rMN7m0ApCowcoDlypBHlkNbp5eJQf/+
|
|||
github.com/prometheus/tsdb v0.7.1/go.mod h1:qhTCs0VvXwvX/y3TZrWD7rabWM+ijKTux40TwIPHuXU=
|
||||
github.com/rakyll/statik v0.1.7 h1:OF3QCZUuyPxuGEP7B4ypUa7sB/iHtqOTDYZXGM8KOdQ=
|
||||
github.com/rakyll/statik v0.1.7/go.mod h1:AlZONWzMtEnMs7W4e/1LURLiI49pIMmp6V9Unghqrcc=
|
||||
github.com/remyoudompheng/bigfft v0.0.0-20190728182440-6a916e37a237 h1:HQagqIiBmr8YXawX/le3+O26N+vPPC1PtjaF3mwnook=
|
||||
github.com/remyoudompheng/bigfft v0.0.0-20190728182440-6a916e37a237/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo=
|
||||
github.com/remyoudompheng/bigfft v0.0.0-20200410134404-eec4a21b6bb0 h1:OdAsTTz6OkFY5QxjkYwrChwuRruF69c169dPK26NUlk=
|
||||
github.com/remyoudompheng/bigfft v0.0.0-20200410134404-eec4a21b6bb0/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo=
|
||||
github.com/rogpeppe/fastuuid v0.0.0-20150106093220-6724a57986af/go.mod h1:XWv6SoW27p1b0cqNHllgS5HIMJraePCO15w5zCzIWYg=
|
||||
github.com/rogpeppe/go-internal v1.3.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4=
|
||||
github.com/rs/cors v1.7.0 h1:+88SsELBHx5r+hZ8TCkggzSstaWNbDvThkVK8H6f9ik=
|
||||
|
|
@ -304,7 +272,6 @@ github.com/russross/blackfriday/v2 v2.0.1/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQD
|
|||
github.com/ryanuber/columnize v0.0.0-20160712163229-9b3edd62028f/go.mod h1:sm1tb6uqfes/u+d4ooFouqFdy9/2g9QGwK3SQygK0Ts=
|
||||
github.com/satori/go.uuid v1.2.0 h1:0uYX9dsZ2yD7q2RtLRtPSdGDWzjeM3TbMJP9utgA0ww=
|
||||
github.com/satori/go.uuid v1.2.0/go.mod h1:dA0hQrYB0VpLJoorglMZABFdXlWrHn1NEOzdhQKdks0=
|
||||
github.com/sean-/seed v0.0.0-20170313163322-e2103e2c3529 h1:nn5Wsu0esKSJiIVhscUtVbo7ada43DJhG55ua/hjS5I=
|
||||
github.com/sean-/seed v0.0.0-20170313163322-e2103e2c3529/go.mod h1:DxrIzT+xaE7yg65j358z/aeFdxmN0P9QXhEzd20vsDc=
|
||||
github.com/shirou/gopsutil/v3 v3.20.11 h1:NeVf1K0cgxsWz+N3671ojRptdgzvp7BXL3KV21R0JnA=
|
||||
github.com/shirou/gopsutil/v3 v3.20.11/go.mod h1:igHnfak0qnw1biGeI2qKQvu0ZkwvEkUcCLlYhZzdr/4=
|
||||
|
|
@ -338,11 +305,9 @@ github.com/spf13/viper v1.7.0/go.mod h1:8WkrPz2fc9jxqZNCJI/76HCieCp4Q8HaLFoCha5q
|
|||
github.com/spf13/viper v1.7.1 h1:pM5oEahlgWv/WnHXpgbKz7iLIxRf65tye2Ci+XFK5sk=
|
||||
github.com/spf13/viper v1.7.1/go.mod h1:8WkrPz2fc9jxqZNCJI/76HCieCp4Q8HaLFoCha5qpdg=
|
||||
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
github.com/stretchr/objx v0.1.1 h1:2vfRuCMp5sSVIDSqO8oNnWJq7mPa6KVP3iPIwFBuy8A=
|
||||
github.com/stretchr/objx v0.1.1/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs=
|
||||
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
|
||||
github.com/stretchr/testify v1.4.0 h1:2E4SXV/wtOkTonXsotYi4li6zVWxYlZuYNCXe9XRJyk=
|
||||
github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4=
|
||||
github.com/stretchr/testify v1.6.1 h1:hDPOHmpOpP40lSULcqw7IrRb/u7w6RpDC9399XyoNd0=
|
||||
github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
|
||||
|
|
@ -456,11 +421,10 @@ golang.org/x/sys v0.0.0-20191001151750-bb3f8db39f24/go.mod h1:h1NjWce9XRLGQEsW7w
|
|||
golang.org/x/sys v0.0.0-20191005200804-aed5e4c7ecf9/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20191220142924-d4481acd189f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20200202164722-d101bd2416d5/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd h1:xhmwyvizuTgC2qz7ZlMluP20uW+C3Rm0FD/WLDX8884=
|
||||
golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20201024232916-9f70ab9862d5/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20201214095126-aec9a390925b h1:tv7/y4pd+sR8bcNb2D6o7BNU6zjWm0VjQLac+w7fNNM=
|
||||
golang.org/x/sys v0.0.0-20201214095126-aec9a390925b/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20210124154548-22da62e12c0c h1:VwygUrnw9jn88c4u8GD3rZQbqrP/tgas88tPUbBxQrk=
|
||||
golang.org/x/sys v0.0.0-20210124154548-22da62e12c0c/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
|
||||
golang.org/x/text v0.3.1-0.20180807135948-17ff2d5776d2/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
|
||||
golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk=
|
||||
|
|
@ -510,7 +474,6 @@ google.golang.org/genproto v0.0.0-20190418145605-e7d98fc518a7/go.mod h1:VzzqZJRn
|
|||
google.golang.org/genproto v0.0.0-20190425155659-357c62f0e4bb/go.mod h1:VzzqZJRnGkLBvHegQrXjBqPurQTc5/KpmUdxsrq26oE=
|
||||
google.golang.org/genproto v0.0.0-20190502173448-54afdca5d873/go.mod h1:VzzqZJRnGkLBvHegQrXjBqPurQTc5/KpmUdxsrq26oE=
|
||||
google.golang.org/genproto v0.0.0-20190801165951-fa694d86fc64/go.mod h1:DMBHOl98Agz4BDEuKkezgsaosCRResVns1a3J2ZsMNc=
|
||||
google.golang.org/genproto v0.0.0-20190819201941-24fa4b261c55 h1:gSJIx1SDwno+2ElGhA4+qG2zF97qiUzTM+rQ0klBOcE=
|
||||
google.golang.org/genproto v0.0.0-20190819201941-24fa4b261c55/go.mod h1:DMBHOl98Agz4BDEuKkezgsaosCRResVns1a3J2ZsMNc=
|
||||
google.golang.org/genproto v0.0.0-20190911173649-1774047e7e51/go.mod h1:IbNlFCBrqXvoKpeg0TB2l7cyZUmoaFKYIwrEpbDKLA8=
|
||||
google.golang.org/genproto v0.0.0-20191108220845-16a3f7862a1a h1:Ob5/580gVHBJZgXnff1cZDbG+xLtMVE5mDRTe+nIsX4=
|
||||
|
|
@ -523,13 +486,6 @@ google.golang.org/grpc v1.25.1/go.mod h1:c3i+UQWmh7LiEpx4sFZnkU36qjEYZ0imhYfXVyQ
|
|||
google.golang.org/grpc v1.26.0/go.mod h1:qbnxyOmOxrQa7FizSgH+ReBfzJrCY1pSN7KXBS8abTk=
|
||||
google.golang.org/grpc v1.28.0 h1:bO/TA4OxCOummhSf10siHuG7vJOiwh7SpRpFZDkOgl4=
|
||||
google.golang.org/grpc v1.28.0/go.mod h1:rpkK4SK4GF4Ach/+MFLZUBavHOvF2JJB5uozKKal+60=
|
||||
google.golang.org/protobuf v0.0.0-20200109180630-ec00e32a8dfd/go.mod h1:DFci5gLYBciE7Vtevhsrf46CRTquxDuWsQurQQe4oz8=
|
||||
google.golang.org/protobuf v0.0.0-20200221191635-4d8936d0db64/go.mod h1:kwYJMbMJ01Woi6D6+Kah6886xMZcty6N08ah7+eCXa0=
|
||||
google.golang.org/protobuf v0.0.0-20200228230310-ab0ca4ff8a60/go.mod h1:cfTl7dwQJ+fmap5saPgwCLgHXTUD7jkjRqWcaiX5VyM=
|
||||
google.golang.org/protobuf v1.20.1-0.20200309200217-e05f789c0967/go.mod h1:A+miEFZTKqfCUM6K7xSMQL9OKL/b6hQv+e19PK+JZNE=
|
||||
google.golang.org/protobuf v1.21.0/go.mod h1:47Nbq4nVaFHyn7ilMalzfO3qCViNmqZ2kzikPIcrTAo=
|
||||
google.golang.org/protobuf v1.23.0 h1:4MY060fB1DLGMB/7MBTLnwQUY6+F09GEiz6SsrNqyzM=
|
||||
google.golang.org/protobuf v1.23.0/go.mod h1:EGpADcykh3NcUnDUJcl1+ZksZNG86OlYog2l/sGQquU=
|
||||
gopkg.in/alecthomas/kingpin.v2 v2.2.6/go.mod h1:FMv+mEhP44yOT+4EoQTLFTRgOQ1FBLkstjWtayDeSgw=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
|
|
@ -542,7 +498,6 @@ gopkg.in/ini.v1 v1.51.0/go.mod h1:pNLf8WUiyNEtQjuu5G5vTm06TEv9tsIgeAvK8hOrP4k=
|
|||
gopkg.in/resty.v1 v1.12.0/go.mod h1:mDo4pnntr5jdWRML875a/NmxYqAlA73dVijT2AXvQQo=
|
||||
gopkg.in/yaml.v2 v2.0.0-20170812160011-eb3733d160e7/go.mod h1:JAlM8MvJe8wmxCU4Bli9HhUf9+ttbYbLASfIpnQbh74=
|
||||
gopkg.in/yaml.v2 v2.2.1/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
||||
gopkg.in/yaml.v2 v2.2.2 h1:ZCJp+EgiOT7lHqUV2J862kp8Qj64Jo6az82+3Td9dZw=
|
||||
gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
||||
gopkg.in/yaml.v2 v2.2.4/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
||||
gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
||||
|
|
@ -555,8 +510,8 @@ honnef.co/go/tools v0.0.0-20190106161140-3f1c8253044a/go.mod h1:rf3lG4BRIbNafJWh
|
|||
honnef.co/go/tools v0.0.0-20190418001031-e561f6794a2a/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4=
|
||||
honnef.co/go/tools v0.0.0-20190523083050-ea95bdfd59fc/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4=
|
||||
honnef.co/go/tools v0.0.1-2019.2.3/go.mod h1:a3bituU0lyd329TUQxRnasdCoJDkEUEAqEt0JzvZhAg=
|
||||
modernc.org/mathutil v1.0.0 h1:93vKjrJopTPrtTNpZ8XIovER7iCIH1QU7wNbOQXC60I=
|
||||
modernc.org/mathutil v1.0.0/go.mod h1:wU0vUrJsVWBZ4P6e7xtFJEhFSNsfRLJ8H458uRjg03k=
|
||||
modernc.org/mathutil v1.2.2 h1:+yFk8hBprV+4c0U9GjFtL+dV3N8hOJ8JCituQcMShFY=
|
||||
modernc.org/mathutil v1.2.2/go.mod h1:mZW8CKdRPY1v87qxC/wUdX5O1qDzXMP5TH3wjfpga6E=
|
||||
modernc.org/strutil v1.0.0 h1:XVFtQwFVwc02Wk+0L/Z/zDDXO81r5Lhe6iMKmGX3KhE=
|
||||
modernc.org/strutil v1.0.0/go.mod h1:lstksw84oURvj9y3tn8lGvRxyRC1S2+g5uuIzNfIOBs=
|
||||
rsc.io/binaryregexp v0.2.0/go.mod h1:qTv7/COck+e2FymRvadv62gMdZztPaShugOCi3I+8D8=
|
||||
|
|
|
|||
51
handler.go
51
handler.go
|
|
@ -16,6 +16,7 @@ package pilosa
|
|||
|
||||
import (
|
||||
"encoding/json"
|
||||
"sort"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/tracing"
|
||||
"github.com/pkg/errors"
|
||||
|
|
@ -98,6 +99,56 @@ func (resp *QueryResponse) MarshalJSON() ([]byte, error) {
|
|||
})
|
||||
}
|
||||
|
||||
// used in the oracle, orders keys lexigraphically
|
||||
func (resp *QueryResponse) Equals(other *QueryResponse) bool {
|
||||
other.Sort()
|
||||
resp.Sort()
|
||||
|
||||
a, err := resp.MarshalJSON()
|
||||
panicOn(err)
|
||||
b, err := other.MarshalJSON()
|
||||
panicOn(err)
|
||||
return string(a) == string(b)
|
||||
|
||||
}
|
||||
|
||||
// SortIfIsRowIdentifiers forces a certain order of rowkeys so that the sauron
|
||||
// tool gets comparable results
|
||||
func SortIfIsRowIdentifiers(i interface{}) (bool, interface{}) {
|
||||
if m, ok := i.(map[string]interface{}); ok {
|
||||
it, ok := m["keys"]
|
||||
if ok {
|
||||
items := it.([]interface{})
|
||||
if len(items) > 0 {
|
||||
sort.Slice(items, func(p, q int) bool {
|
||||
return items[p].(string) < items[q].(string)
|
||||
})
|
||||
return true, items
|
||||
//return true, ToGenericArray(sitems)
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
return false, nil
|
||||
}
|
||||
|
||||
// TODO (twg) probably going to get fancy here but hardcode to test idea
|
||||
func (resp *QueryResponse) Sort() bool {
|
||||
modified := false
|
||||
if resp == nil || resp.Results == nil {
|
||||
return false
|
||||
}
|
||||
for i, result := range resp.Results {
|
||||
sorted, results := SortIfIsRowIdentifiers(result)
|
||||
if sorted {
|
||||
resp.Results[i] = results
|
||||
modified = true
|
||||
}
|
||||
}
|
||||
|
||||
return modified
|
||||
}
|
||||
|
||||
// Handler is the interface for the data handler, a wrapper around
|
||||
// Pilosa's data store.
|
||||
type Handler interface {
|
||||
|
|
|
|||
232
holder.go
232
holder.go
|
|
@ -906,13 +906,13 @@ func (h *Holder) limitedSchema() ([]*IndexInfo, error) {
|
|||
}
|
||||
|
||||
func (h *Holder) schema(ctx context.Context, includeViews bool) ([]*IndexInfo, error) {
|
||||
var a []*IndexInfo
|
||||
|
||||
schema, err := h.schemator.Schema(ctx)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting schema via schemator")
|
||||
}
|
||||
|
||||
a := make([]*IndexInfo, 0, len(schema))
|
||||
|
||||
for _, index := range schema {
|
||||
cim, err := decodeCreateIndexMessage(h.serializer, index.Data)
|
||||
if err != nil {
|
||||
|
|
@ -1478,6 +1478,10 @@ func (h *Holder) logStartup() error {
|
|||
time := time.Now().Format(RFC3339NanoFixedWidth)
|
||||
logLine := fmt.Sprintf("%s\t%s\n", time, Version)
|
||||
|
||||
if err := os.MkdirAll(h.path, 0777); err != nil {
|
||||
return errors.Wrap(err, "creating data directory")
|
||||
}
|
||||
|
||||
f, err := os.OpenFile(h.path+"/startup.log", os.O_APPEND|os.O_WRONLY|os.O_CREATE, 0600)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "opening startup log")
|
||||
|
|
@ -1743,16 +1747,11 @@ func (s *holderSyncer) resetTranslationSync() error {
|
|||
// Set read-only flag for all translation stores.
|
||||
s.setTranslateReadOnlyFlags(snap)
|
||||
|
||||
// Connect to each node that has a primary for which we are a replica.
|
||||
if err := s.initializeIndexTranslateReplication(snap); err != nil {
|
||||
return errors.Wrap(err, "initialize index translate replication")
|
||||
}
|
||||
|
||||
// Connect to primary to stream field data.
|
||||
if err := s.initializeFieldTranslateReplication(snap); err != nil {
|
||||
return errors.Wrap(err, "initialize field translate replication")
|
||||
if err := s.initializeReplication(snap); err != nil {
|
||||
return errors.Wrap(err, "initializing translation replication")
|
||||
}
|
||||
return nil
|
||||
|
||||
}
|
||||
|
||||
////////////////////////////////////////////////////////////
|
||||
|
|
@ -1859,11 +1858,83 @@ func (s *holderSyncer) setTranslateReadOnlyFlags(snap *topology.ClusterSnapshot)
|
|||
s.Cluster.mu.RUnlock()
|
||||
}
|
||||
|
||||
// initializeIndexTranslateReplication connects to each node that is the
|
||||
// primary for a partition that we are a replica of.
|
||||
func (s *holderSyncer) initializeIndexTranslateReplication(snap *topology.ClusterSnapshot) error {
|
||||
// initializeReplication builds a map of nodes for which we need to replicate
|
||||
// any key translation, whether that's field keys (every node replicates these
|
||||
// from the primary) or index keys (only the replica nodes for each partition
|
||||
// replicate these from whichever node is primary for that partition).
|
||||
func (s *holderSyncer) initializeReplication(snap *topology.ClusterSnapshot) error {
|
||||
nodeMaps := make(map[string]TranslateOffsetMap)
|
||||
|
||||
if snap.ReplicaN > 1 {
|
||||
if err := s.populateIndexReplication(nodeMaps, snap); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if err := s.populateFieldReplication(nodeMaps, snap); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for _, node := range snap.Nodes {
|
||||
m := nodeMaps[node.ID]
|
||||
if m.Empty() {
|
||||
continue
|
||||
}
|
||||
|
||||
// Connect to remote node and begin streaming.
|
||||
rd, err := s.Holder.OpenTranslateReader(context.Background(), node.URI.String(), m)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
s.readers = append(s.readers, rd)
|
||||
|
||||
s.syncers.Go(func() error {
|
||||
defer rd.Close()
|
||||
s.readBothTranslateReader(rd, snap)
|
||||
return nil
|
||||
})
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// populateFieldReplication populates a map from node IDs to TranslateOffsetMaps
|
||||
// to record that we need to translate fields which have key translation
|
||||
// from the primary node.
|
||||
func (s *holderSyncer) populateFieldReplication(nodeMaps map[string]TranslateOffsetMap, snap *topology.ClusterSnapshot) error {
|
||||
// Set up field translation
|
||||
if !snap.IsPrimaryFieldTranslationNode(s.Cluster.Node.ID) {
|
||||
primaryID := snap.PrimaryFieldTranslationNode().ID
|
||||
// Build a map of field key offsets to stream from.
|
||||
m := nodeMaps[primaryID]
|
||||
if m == nil {
|
||||
m = make(TranslateOffsetMap)
|
||||
nodeMaps[primaryID] = m
|
||||
}
|
||||
for _, index := range s.Holder.Indexes() {
|
||||
for _, field := range index.Fields() {
|
||||
store := field.TranslateStore()
|
||||
// I think right now this is supposed to be impossible;
|
||||
// we use an InMemTranslateStore by default even if
|
||||
// no translate store is being used or attempted.
|
||||
if store == nil {
|
||||
return fmt.Errorf("no translate store for field %q/%q", index.Name(), field.Name())
|
||||
}
|
||||
offset, err := store.MaxID()
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "cannot determine max id for %q/%q", index.Name(), field.Name())
|
||||
}
|
||||
m.SetFieldOffset(index.Name(), field.Name(), offset)
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// populateIndexReplication populates a map of node IDs to TranslateOffsetMaps
|
||||
// to record which nodes we need to replicate index key translation for.
|
||||
// That means nodes which are the primary for a partition that we're a
|
||||
// non-primary replica for.
|
||||
func (s *holderSyncer) populateIndexReplication(nodeMaps map[string]TranslateOffsetMap, snap *topology.ClusterSnapshot) error {
|
||||
for _, node := range snap.Nodes {
|
||||
// Skip local node.
|
||||
if node.ID == s.Node.ID {
|
||||
continue
|
||||
}
|
||||
|
|
@ -1883,6 +1954,9 @@ func (s *holderSyncer) initializeIndexTranslateReplication(snap *topology.Cluste
|
|||
}
|
||||
|
||||
store := index.TranslateStore(partitionID)
|
||||
if store == nil {
|
||||
return fmt.Errorf("no store available for index %q, partition %d", index.Name(), partitionID)
|
||||
}
|
||||
offset, err := store.MaxID()
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "cannot determine max id for %q", index.Name())
|
||||
|
|
@ -1895,109 +1969,49 @@ func (s *holderSyncer) initializeIndexTranslateReplication(snap *topology.Cluste
|
|||
if len(m) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
// Connect to remote node and begin streaming.
|
||||
rd, err := s.Holder.OpenTranslateReader(context.Background(), node.URI.String(), m)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
s.readers = append(s.readers, rd)
|
||||
|
||||
s.syncers.Go(func() error {
|
||||
defer rd.Close()
|
||||
s.readIndexTranslateReader(rd)
|
||||
return nil
|
||||
})
|
||||
nodeMaps[node.ID] = m
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// initializeFieldTranslateReplication connects the primary to stream field data.
|
||||
func (s *holderSyncer) initializeFieldTranslateReplication(snap *topology.ClusterSnapshot) error {
|
||||
// Skip if primary.
|
||||
if snap.IsPrimaryFieldTranslationNode(s.Cluster.Node.ID) {
|
||||
return nil
|
||||
}
|
||||
// readBothTranslateReader reads key translation for field keys or
|
||||
// index keys from a remote node. Both field and index keys may be sent,
|
||||
// the distinction is that field keys have a non-empty field name.
|
||||
func (s *holderSyncer) readBothTranslateReader(rd TranslateEntryReader, snap *topology.ClusterSnapshot) {
|
||||
for {
|
||||
var entry TranslateEntry
|
||||
if err := rd.ReadEntry(&entry); err != nil {
|
||||
s.Holder.Logger.Printf("cannot read translate entry: %s", err)
|
||||
return
|
||||
}
|
||||
|
||||
// Build a map of partition offsets to stream from.
|
||||
m := make(TranslateOffsetMap)
|
||||
for _, index := range s.Holder.Indexes() {
|
||||
for _, field := range index.Fields() {
|
||||
store := field.TranslateStore()
|
||||
offset, err := store.MaxID()
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "cannot determine max id for %q/%q", index.Name(), field.Name())
|
||||
var store TranslateStore
|
||||
if entry.Field != "" {
|
||||
// Find appropriate store.
|
||||
f := s.Holder.Field(entry.Index, entry.Field)
|
||||
if f == nil {
|
||||
s.Holder.Logger.Printf("field not found: %s/%s", entry.Index, entry.Field)
|
||||
return
|
||||
}
|
||||
store = f.TranslateStore()
|
||||
if store == nil {
|
||||
s.Holder.Logger.Printf("no translate store suitable for index %q, field %q, key %q", entry.Index, entry.Field, entry.Key)
|
||||
return
|
||||
}
|
||||
} else {
|
||||
// Find appropriate store.
|
||||
idx := s.Holder.Index(entry.Index)
|
||||
if idx == nil {
|
||||
s.Holder.Logger.Printf("index not found: %q", entry.Index)
|
||||
return
|
||||
}
|
||||
store = idx.TranslateStore(snap.KeyToKeyPartition(entry.Index, entry.Key))
|
||||
if store == nil {
|
||||
s.Holder.Logger.Printf("no translate store suitable for index %q, key %q", entry.Index, entry.Key)
|
||||
return
|
||||
}
|
||||
m.SetFieldOffset(index.Name(), field.Name(), offset)
|
||||
}
|
||||
}
|
||||
|
||||
// Skip if no replication required.
|
||||
if len(m) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Connect to primary and begin streaming.
|
||||
primary := snap.PrimaryFieldTranslationNode()
|
||||
rd, err := s.Holder.OpenTranslateReader(context.Background(), primary.URI.String(), m)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
s.readers = append(s.readers, rd)
|
||||
|
||||
s.syncers.Go(func() error {
|
||||
defer rd.Close()
|
||||
s.readFieldTranslateReader(rd)
|
||||
return nil
|
||||
})
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *holderSyncer) readIndexTranslateReader(rd TranslateEntryReader) {
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN)
|
||||
|
||||
for {
|
||||
var entry TranslateEntry
|
||||
if err := rd.ReadEntry(&entry); err != nil {
|
||||
s.Holder.Logger.Printf("cannot read index translate entry: %s", err)
|
||||
return
|
||||
}
|
||||
|
||||
// Find appropriate store.
|
||||
idx := s.Holder.Index(entry.Index)
|
||||
if idx == nil {
|
||||
s.Holder.Logger.Printf("index not found: %q", entry.Index)
|
||||
return
|
||||
}
|
||||
|
||||
// Apply replication to store.
|
||||
store := idx.TranslateStore(snap.KeyToKeyPartition(entry.Index, entry.Key))
|
||||
if err := store.ForceSet(entry.ID, entry.Key); err != nil {
|
||||
s.Holder.Logger.Printf("cannot force set index translation data: %d=%q", entry.ID, entry.Key)
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *holderSyncer) readFieldTranslateReader(rd TranslateEntryReader) {
|
||||
for {
|
||||
var entry TranslateEntry
|
||||
if err := rd.ReadEntry(&entry); err != nil {
|
||||
s.Holder.Logger.Printf("cannot read field translate entry: %s", err)
|
||||
return
|
||||
}
|
||||
|
||||
// Find appropriate store.
|
||||
f := s.Holder.Field(entry.Index, entry.Field)
|
||||
if f == nil {
|
||||
s.Holder.Logger.Printf("field not found: %s/%s", entry.Index, entry.Field)
|
||||
return
|
||||
}
|
||||
|
||||
// Apply replication to store.
|
||||
store := f.TranslateStore()
|
||||
if err := store.ForceSet(entry.ID, entry.Key); err != nil {
|
||||
s.Holder.Logger.Printf("cannot force set field translation data: %d=%q", entry.ID, entry.Key)
|
||||
return
|
||||
|
|
|
|||
|
|
@ -22,6 +22,7 @@ import (
|
|||
"io"
|
||||
"io/ioutil"
|
||||
"net/http"
|
||||
"reflect"
|
||||
"sync"
|
||||
|
||||
"github.com/pilosa/pilosa/v2"
|
||||
|
|
@ -33,8 +34,12 @@ func GetOpenTranslateReaderFunc(client *http.Client) pilosa.OpenTranslateReaderF
|
|||
}
|
||||
|
||||
func GetOpenTranslateReaderWithLockerFunc(client *http.Client, locker sync.Locker) pilosa.OpenTranslateReaderFunc {
|
||||
lockType := reflect.TypeOf(locker)
|
||||
if lockType.Kind() == reflect.Ptr {
|
||||
lockType = lockType.Elem()
|
||||
}
|
||||
return func(ctx context.Context, nodeURL string, offsets pilosa.TranslateOffsetMap) (pilosa.TranslateEntryReader, error) {
|
||||
return openTranslateReader(ctx, nodeURL, offsets, client, locker)
|
||||
return openTranslateReader(ctx, nodeURL, offsets, client, reflect.New(lockType).Interface().(sync.Locker))
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -17,18 +17,17 @@ package pilosa
|
|||
import (
|
||||
"crypto/rand"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"reflect"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
bolt "go.etcd.io/bbolt"
|
||||
)
|
||||
|
||||
func TestIDAlloc(t *testing.T) {
|
||||
// Acquire a temporary file.
|
||||
f, err := ioutil.TempFile("", "")
|
||||
f, err := testhook.TempFile(t, "idalloc")
|
||||
if err != nil {
|
||||
t.Errorf("acquiring temporary file: %v", err)
|
||||
return
|
||||
|
|
@ -42,10 +41,6 @@ func TestIDAlloc(t *testing.T) {
|
|||
|
||||
// Open bolt.
|
||||
db, err := bolt.Open(f.Name(), 0666, &bolt.Options{Timeout: 1 * time.Second})
|
||||
if rerr := os.Remove(f.Name()); rerr != nil {
|
||||
t.Errorf("removing temporary file: %v", rerr)
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
t.Errorf("opening bolt: %v", err)
|
||||
return
|
||||
|
|
|
|||
|
|
@ -29,6 +29,8 @@ import (
|
|||
"io/ioutil"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
)
|
||||
|
||||
// TestReopenAppend -- make sure we always append to an existing file
|
||||
|
|
@ -42,16 +44,12 @@ import (
|
|||
// 5. read file, make sure it contains line0,line1,line2
|
||||
//
|
||||
func TestReopenAppend(t *testing.T) {
|
||||
// TODO fix
|
||||
// (travis) I have no idea what this TODO is asking for.
|
||||
// Perhaps use `ioutil.TempFile()`?
|
||||
var fname = "/tmp/foo"
|
||||
|
||||
// Step 1 -- Create a sample file using normal means
|
||||
forig, err := os.Create(fname)
|
||||
forig, err := testhook.TempFile(t, "logger-reopen")
|
||||
if err != nil {
|
||||
t.Fatalf("Unable to create initial file %s: %s", fname, err)
|
||||
t.Fatalf("unable to create initial file: %v", err)
|
||||
}
|
||||
fname := forig.Name()
|
||||
|
||||
_, err = forig.Write([]byte("line0\n"))
|
||||
if err != nil {
|
||||
t.Fatalf("Unable to write initial line %s: %s", fname, err)
|
||||
|
|
@ -96,26 +94,23 @@ func TestReopenAppend(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
// Test that reopen works when Inode is swapped out
|
||||
// 1. Create a sample file using normal means
|
||||
// 2. Open a ioreopen.File
|
||||
// write line 1
|
||||
// 3. call Reopen
|
||||
// write line 2
|
||||
// 4. close file
|
||||
// 5. read file, make sure it contains line0,line1,line2
|
||||
// Test that reopen works when Inode is swapped out. That is to say,
|
||||
// if the previous file has been renamed, we want to get a new file
|
||||
// that isn't the old one, not keep writing to the old file.
|
||||
//
|
||||
// 1. Create a sample file using normal means
|
||||
// 2. Write to it.
|
||||
// 3. Rename it.
|
||||
// 4. Call reopen.
|
||||
// 5. Write line 2.
|
||||
// 6. Read file, expecting to see only line 2.
|
||||
func TestChangeInode(t *testing.T) {
|
||||
// TODO fix
|
||||
// (travis) I have no idea what this TODO is asking for.
|
||||
// Perhaps use `ioutil.TempFile()`?
|
||||
var fname = "/tmp/foo"
|
||||
|
||||
// Step 1 -- Create a empty sample file
|
||||
forig, err := os.Create(fname)
|
||||
forig, err := testhook.TempFile(t, "changeInode")
|
||||
if err != nil {
|
||||
t.Fatalf("Unable to create initial file %s: %s", fname, err)
|
||||
t.Fatalf("Unable to create initial file: %s", err)
|
||||
}
|
||||
fname := forig.Name()
|
||||
err = forig.Close()
|
||||
if err != nil {
|
||||
t.Fatalf("Unable to close initial file: %s", err)
|
||||
|
|
@ -136,6 +131,8 @@ func TestChangeInode(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Errorf("Renaming error: %s", err)
|
||||
}
|
||||
// remove the scratch file
|
||||
defer os.Remove(fname + ".orig")
|
||||
_, err = f.Write([]byte("after1\n"))
|
||||
if err != nil {
|
||||
t.Errorf("Write error: %s", err)
|
||||
|
|
|
|||
24
rbf.go
24
rbf.go
|
|
@ -284,15 +284,21 @@ func (tx *RBFTx) addOrRemove(index, field, view string, shard uint64, batched, r
|
|||
// not first time through, write what we got.
|
||||
if remove && (rc == nil || rc.N() == 0) {
|
||||
err = tx.RemoveContainer(index, field, view, shard, lastHi)
|
||||
panicOn(err)
|
||||
if err != nil {
|
||||
return 0, errors.Wrap(err, "failed to remove container")
|
||||
}
|
||||
} else {
|
||||
err = tx.PutContainer(index, field, view, shard, lastHi, rc)
|
||||
panicOn(err)
|
||||
if err != nil {
|
||||
return 0, errors.Wrap(err, "failed to put container")
|
||||
}
|
||||
}
|
||||
}
|
||||
// get the next container
|
||||
rc, err = tx.Container(index, field, view, shard, hi)
|
||||
panicOn(err)
|
||||
if err != nil {
|
||||
return 0, errors.Wrap(err, "failed to retrieve container")
|
||||
}
|
||||
} // else same container, keep adding bits to rct.
|
||||
chng := false
|
||||
// rc can be nil before, and nil after, in both Remove/Add below.
|
||||
|
|
@ -311,17 +317,23 @@ func (tx *RBFTx) addOrRemove(index, field, view string, shard uint64, batched, r
|
|||
if remove {
|
||||
if rc == nil || rc.N() == 0 {
|
||||
err = tx.RemoveContainer(index, field, view, shard, hi)
|
||||
panicOn(err)
|
||||
if err != nil {
|
||||
return 0, errors.Wrap(err, "failed to remove container")
|
||||
}
|
||||
} else {
|
||||
err = tx.PutContainer(index, field, view, shard, hi, rc)
|
||||
panicOn(err)
|
||||
if err != nil {
|
||||
return 0, errors.Wrap(err, "failed to put container")
|
||||
}
|
||||
}
|
||||
} else {
|
||||
if rc == nil || rc.N() == 0 {
|
||||
panic("there should be no way to have an empty bitmap AFTER an Add() operation")
|
||||
}
|
||||
err = tx.PutContainer(index, field, view, shard, hi, rc)
|
||||
panicOn(err)
|
||||
if err != nil {
|
||||
return 0, errors.Wrap(err, "failed to put container")
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
|
|
|||
|
|
@ -32,7 +32,7 @@ import (
|
|||
)
|
||||
|
||||
func TestDB_Open(t *testing.T) {
|
||||
db := NewDB()
|
||||
db := NewDB(t)
|
||||
if err := db.Open(); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if err := db.Close(); err != nil {
|
||||
|
|
|
|||
|
|
@ -16,17 +16,19 @@ package rbf
|
|||
|
||||
import (
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
"runtime/pprof"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
//"time"
|
||||
|
||||
"github.com/pilosa/pilosa/v2/rbf/cfg"
|
||||
"github.com/pilosa/pilosa/v2/roaring"
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
|
||||
// "github.com/pilosa/pilosa/v2/txkey"
|
||||
txkey "github.com/pilosa/pilosa/v2/short_txkey"
|
||||
)
|
||||
|
|
@ -71,7 +73,7 @@ func TestIngest_lots_of_views(t *testing.T) {
|
|||
//vv("m1.TotalAlloc = %v", m1.TotalAlloc)
|
||||
}()
|
||||
|
||||
path, err := ioutil.TempDir("", "rbf_ingest_lots_of_views")
|
||||
path, err := testhook.TempDir(t, "rbf_ingest_lots_of_views")
|
||||
panicOn(err)
|
||||
defer os.Remove(path)
|
||||
|
||||
|
|
|
|||
|
|
@ -18,7 +18,6 @@ import (
|
|||
"flag"
|
||||
"fmt"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"math/rand"
|
||||
"os"
|
||||
"runtime"
|
||||
|
|
@ -27,6 +26,7 @@ import (
|
|||
|
||||
"github.com/pilosa/pilosa/v2/rbf"
|
||||
rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg"
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
)
|
||||
|
||||
var quickCheckN *int = flag.Int("quickchecks", 10, "The number of iterations for each quickcheck")
|
||||
|
|
@ -61,8 +61,8 @@ func TestReadWriteRootRecord(t *testing.T) {
|
|||
}
|
||||
|
||||
// NewDB returns a new instance of DB with a temporary path.
|
||||
func NewDB(cfg ...*rbfcfg.Config) *rbf.DB {
|
||||
path, err := ioutil.TempDir("", "")
|
||||
func NewDB(tb testing.TB, cfg ...*rbfcfg.Config) *rbf.DB {
|
||||
path, err := testhook.TempDir(tb, "rbfdb")
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
|
@ -78,7 +78,7 @@ func NewDB(cfg ...*rbfcfg.Config) *rbf.DB {
|
|||
// MustOpenDB returns a db opened on a temporary file. On error, fail test.
|
||||
func MustOpenDB(tb testing.TB, cfg ...*rbfcfg.Config) *rbf.DB {
|
||||
tb.Helper()
|
||||
db := NewDB(cfg...)
|
||||
db := NewDB(tb, cfg...)
|
||||
if err := db.Open(); err != nil {
|
||||
tb.Fatal(err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -16,12 +16,12 @@ package rbf
|
|||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg"
|
||||
"github.com/pilosa/pilosa/v2/roaring"
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
)
|
||||
|
||||
// util_test adds reusable utilities for testing.
|
||||
|
|
@ -144,7 +144,7 @@ func verifyElemNBitN(tx *Tx, lc leafCell) {
|
|||
func testHelperMustOpenNewDB(tb testing.TB, cfg ...*rbfcfg.Config) *DB {
|
||||
tb.Helper()
|
||||
|
||||
path, err := ioutil.TempDir("", "")
|
||||
path, err := testhook.TempDir(tb, "rbfdb")
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -192,7 +192,7 @@ func TestClusterResize_AddNode(t *testing.T) {
|
|||
}
|
||||
lsns[i] = l.(*net.TCPListener)
|
||||
}
|
||||
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
|
||||
portsCfg := test.GenPortsConfig(t, test.NewPorts(lsns))
|
||||
|
||||
m1.Config.Etcd = portsCfg[0].Etcd
|
||||
m1.Config.Name = portsCfg[0].Name
|
||||
|
|
|
|||
|
|
@ -106,7 +106,7 @@ func TestHandler_Endpoints(t *testing.T) {
|
|||
t.Fatalf("unexpected status code: %d", w.Code)
|
||||
}
|
||||
body := w.Body.String()
|
||||
if body != "{\"indexes\":null}\n" {
|
||||
if body != "{\"indexes\":[]}\n" {
|
||||
t.Fatalf("unexpected empty schema: '%v'", body)
|
||||
}
|
||||
|
||||
|
|
@ -119,7 +119,7 @@ func TestHandler_Endpoints(t *testing.T) {
|
|||
t.Fatalf("unexpected status code: %d", w.Code)
|
||||
}
|
||||
body := w.Body.String()
|
||||
if body != "{\"indexes\":null}\n" {
|
||||
if body != "{\"indexes\":[]}\n" {
|
||||
t.Fatalf("unexpected empty schema: '%v'", body)
|
||||
}
|
||||
|
||||
|
|
@ -790,7 +790,7 @@ func TestHandler_Endpoints(t *testing.T) {
|
|||
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i0/query", strings.NewReader(`TopN(f0, n=2)`)))
|
||||
if w.Code != gohttp.StatusOK {
|
||||
t.Fatalf("unexpected status code: %d", w.Code)
|
||||
} else if body := w.Body.String(); body != `{"results":[[{"id":30,"key":"","count":3},{"id":31,"key":"","count":1}]]}`+"\n" {
|
||||
} else if body := w.Body.String(); body != `{"results":[[{"count":3,"id":30,"key":""},{"count":1,"id":31,"key":""}]]}`+"\n" {
|
||||
t.Fatalf("unexpected body: %q", body)
|
||||
}
|
||||
})
|
||||
|
|
|
|||
|
|
@ -16,7 +16,6 @@ package test
|
|||
|
||||
import (
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"net"
|
||||
"strings"
|
||||
"testing"
|
||||
|
|
@ -121,7 +120,7 @@ func GetPortsGenConfigs(tb testing.TB, nodes []*Command) error {
|
|||
}
|
||||
|
||||
//GenPortsConfig creates specific configuration for etcd.
|
||||
func GenPortsConfig(ports []Ports) []*server.Config {
|
||||
func GenPortsConfig(tb testing.TB, ports []Ports) []*server.Config {
|
||||
cfgs := make([]*server.Config, len(ports))
|
||||
clusterURLs := make([]string, len(ports))
|
||||
for i := range cfgs {
|
||||
|
|
@ -134,7 +133,7 @@ func GenPortsConfig(ports []Ports) []*server.Config {
|
|||
lPeerURL := fmt.Sprintf("http://localhost:%d", portP)
|
||||
|
||||
discoDir := ""
|
||||
if d, err := ioutil.TempDir("", "disco."); err == nil {
|
||||
if d, err := testhook.TempDir(tb, "disco."); err == nil {
|
||||
discoDir = d
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -87,12 +87,27 @@ func TempDir(tb testing.TB, pattern string) (path string, err error) {
|
|||
if err == nil {
|
||||
Cleanup(tb, func() {
|
||||
os.RemoveAll(path)
|
||||
fmt.Println("--- testhook:", path, tb.Name())
|
||||
// fmt.Println("--- testhook: cleaning up dir", path, tb.Name())
|
||||
})
|
||||
}
|
||||
return path, err
|
||||
}
|
||||
|
||||
// TempFile creates a temp file that will be automatically deleted when
|
||||
// this test completes, using go1.14's [TB].Cleanup() if available.
|
||||
func TempFile(tb testing.TB, pattern string) (file *os.File, err error) {
|
||||
file, err = ioutil.TempFile("", pattern)
|
||||
if err == nil {
|
||||
path := file.Name()
|
||||
Cleanup(tb, func() {
|
||||
file.Close()
|
||||
os.Remove(path)
|
||||
// fmt.Println("--- testhook: cleaning up file", path, tb.Name())
|
||||
})
|
||||
}
|
||||
return file, err
|
||||
}
|
||||
|
||||
// TempDirInDir creates a temp directory that will be automatically deleted when
|
||||
// this test completes, using go1.14's [TB].Cleanup(), but with a specified
|
||||
// path instead of the default Go TMPDIR. Only some tests use this, which is
|
||||
|
|
@ -106,3 +121,19 @@ func TempDirInDir(tb testing.TB, dir string, pattern string) (path string, err e
|
|||
}
|
||||
return path, err
|
||||
}
|
||||
|
||||
// TempFileInDir creates a temp file that will be automatically deleted when
|
||||
// this test completes, using go1.14's [TB].Cleanup(), but with a specified
|
||||
// path instead of the default Go TMPDIR. Only some tests use this, which is
|
||||
// possibly an error...
|
||||
func TempFileInDir(tb testing.TB, dir string, pattern string) (file *os.File, err error) {
|
||||
file, err = ioutil.TempFile(dir, pattern)
|
||||
if err == nil {
|
||||
path := file.Name()
|
||||
Cleanup(tb, func() {
|
||||
file.Close()
|
||||
os.Remove(path)
|
||||
})
|
||||
}
|
||||
return file, err
|
||||
}
|
||||
|
|
|
|||
17
translate.go
17
translate.go
|
|
@ -304,6 +304,18 @@ func (m TranslateOffsetMap) FieldOffset(index, name string) uint64 {
|
|||
return m[index].Fields[name]
|
||||
}
|
||||
|
||||
// Empty reports whether there are any actual entries in the map. This
|
||||
// is distinct from len(m) == 0 in that an entry in this map which is
|
||||
// itself empty doesn't count as non-empty.
|
||||
func (m TranslateOffsetMap) Empty() bool {
|
||||
for _, sub := range m {
|
||||
if !sub.Empty() {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// SetFieldOffset sets the offset for the given field.
|
||||
func (m TranslateOffsetMap) SetFieldOffset(index, name string, offset uint64) {
|
||||
if m[index] == nil {
|
||||
|
|
@ -317,6 +329,11 @@ type IndexTranslateOffsetMap struct {
|
|||
Fields map[string]uint64 `json:"fields"`
|
||||
}
|
||||
|
||||
// Empty reports whether this map has neither partitions nor fields.
|
||||
func (i *IndexTranslateOffsetMap) Empty() bool {
|
||||
return len(i.Partitions) == 0 && len(i.Fields) == 0
|
||||
}
|
||||
|
||||
func NewIndexTranslateOffsetMap() *IndexTranslateOffsetMap {
|
||||
return &IndexTranslateOffsetMap{
|
||||
Partitions: make(map[int]uint64),
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue