mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-08 03:47:51 +00:00
Merge branch 'master' into 1370-cluster-locking
This commit is contained in:
commit
d8597f320a
37 changed files with 265 additions and 205 deletions
|
|
@ -21,10 +21,8 @@ jobs:
|
|||
include:
|
||||
- stage: metalinter
|
||||
install:
|
||||
- make -B install-dep vendor
|
||||
- go get -u github.com/alecthomas/gometalinter
|
||||
- gometalinter --install
|
||||
script: make metalinter
|
||||
- make -B install-dep vendor install-gometalinter
|
||||
script: make gometalinter
|
||||
env:
|
||||
- GOARCH=amd64
|
||||
- stage: deploy
|
||||
|
|
|
|||
30
Makefile
30
Makefile
|
|
@ -1,4 +1,4 @@
|
|||
.PHONY: build check-clean clean cover cover-viz default docker docker-build docker-test generate generate-protoc generate-pql install install-build-deps install-dep install-protoc install-protoc-gen-gofast install-peg metalinter prerelease prerelease-upload release release-build require-dep require-protoc require-protoc-gen-gofast require-peg test
|
||||
.PHONY: build check-clean clean cover cover-viz default docker docker-build docker-test generate generate-protoc generate-pql gometalinter install install-build-deps install-dep install-gometalinter install-protoc install-protoc-gen-gofast install-peg prerelease prerelease-upload release release-build require-dep require-gometalinter require-protoc require-protoc-gen-gofast require-peg test
|
||||
|
||||
CLONE_URL=github.com/pilosa/pilosa
|
||||
VERSION := $(shell git describe --tags 2> /dev/null || echo unknown)
|
||||
|
|
@ -108,8 +108,24 @@ docker-build:
|
|||
docker-test:
|
||||
docker run --rm -v $(PWD):/go/src/$(CLONE_URL) -w /go/src/$(CLONE_URL) golang:$(GO_VERSION) go test -tags='$(BUILD_TAGS)' $(TESTFLAGS) ./...
|
||||
|
||||
metalinter:
|
||||
gometalinter --vendor --disable-all --enable=gotype --enable=gotypex --enable=gofmt --enable=goimports --enable=interfacer --enable=misspell --enable=unparam --deadline=60s --exclude "^internal/.*\.pb\.go" ./...
|
||||
# Run gometalinter with custom flags
|
||||
gometalinter: require-gometalinter
|
||||
gometalinter --vendor --disable-all \
|
||||
--deadline=60s \
|
||||
--enable=deadcode \
|
||||
--enable=gochecknoinits \
|
||||
--enable=gofmt \
|
||||
--enable=goimports \
|
||||
--enable=gotype \
|
||||
--enable=gotypex \
|
||||
--enable=ineffassign \
|
||||
--enable=interfacer \
|
||||
--enable=misspell \
|
||||
--enable=nakedret \
|
||||
--enable=unparam \
|
||||
--exclude "^internal/.*\.pb\.go" \
|
||||
--exclude "^pql/pql.peg.go" \
|
||||
./...
|
||||
|
||||
######################
|
||||
# Build dependencies #
|
||||
|
|
@ -134,6 +150,9 @@ require-protoc:
|
|||
require-peg:
|
||||
$(call require,peg)
|
||||
|
||||
require-gometalinter:
|
||||
$(call require,gometalinter)
|
||||
|
||||
install-build-deps: install-dep install-protoc-gen-gofast install-protoc install-stringer install-peg
|
||||
|
||||
install-dep:
|
||||
|
|
@ -150,3 +169,8 @@ install-protoc:
|
|||
|
||||
install-peg:
|
||||
go get github.com/pointlander/peg
|
||||
|
||||
install-gometalinter:
|
||||
go get -u github.com/alecthomas/gometalinter
|
||||
gometalinter --install
|
||||
go get github.com/remyoudompheng/go-misc/deadcode
|
||||
|
|
|
|||
6
attr.go
6
attr.go
|
|
@ -204,12 +204,6 @@ func DecodeAttrs(v []byte) (map[string]interface{}, error) {
|
|||
return decodeAttrs(pb.GetAttrs()), nil
|
||||
}
|
||||
|
||||
func newMemAttrStore() AttrStore {
|
||||
return &memAttrStore{
|
||||
store: make(map[uint64]map[string]interface{}),
|
||||
}
|
||||
}
|
||||
|
||||
// memAttrStore represents an in-memory implementation of the AttrStore interface.
|
||||
type memAttrStore struct {
|
||||
store map[uint64]map[string]interface{}
|
||||
|
|
|
|||
|
|
@ -143,7 +143,7 @@ func (s *attrStore) Attrs(id uint64) (m map[string]interface{}, err error) {
|
|||
// Add to cache.
|
||||
s.attrCache.Set(id, m)
|
||||
|
||||
return
|
||||
return m, nil
|
||||
}
|
||||
|
||||
// SetAttrs sets attribute values for a given ID.
|
||||
|
|
|
|||
|
|
@ -37,12 +37,8 @@ type broadcaster interface {
|
|||
// TODO add at least a single "isMessage()" method.
|
||||
type Message interface{}
|
||||
|
||||
func init() {
|
||||
NopBroadcaster = &nopBroadcaster{}
|
||||
}
|
||||
|
||||
// NopBroadcaster represents a Broadcaster that doesn't do anything.
|
||||
var NopBroadcaster broadcaster
|
||||
var NopBroadcaster broadcaster = &nopBroadcaster{}
|
||||
|
||||
type nopBroadcaster struct{}
|
||||
|
||||
|
|
|
|||
|
|
@ -808,9 +808,6 @@ type Hasher interface {
|
|||
Hash(key uint64, n int) int
|
||||
}
|
||||
|
||||
// newHasher returns a new instance of the default hasher.
|
||||
func newHasher() Hasher { return &jmphasher{} }
|
||||
|
||||
// jmphasher represents an implementation of jmphash. Implements Hasher.
|
||||
type jmphasher struct{}
|
||||
|
||||
|
|
|
|||
|
|
@ -372,7 +372,8 @@ func TestHasher(t *testing.T) {
|
|||
{0x0ddc0ffeebadf00d, []int{0, 1, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2, 2, 15, 15, 15, 15}},
|
||||
} {
|
||||
for i, v := range tt.bucket {
|
||||
if got := newHasher().Hash(tt.key, i+1); got != v {
|
||||
hasher := &jmphasher{}
|
||||
if got := hasher.Hash(tt.key, i+1); got != v {
|
||||
t.Errorf("hash(%v,%v)=%v, want %v", tt.key, i+1, got, v)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -27,7 +27,7 @@ import (
|
|||
|
||||
var checker *ctl.CheckCommand
|
||||
|
||||
func newCheckCommand(stdin io.Reader, stdout, stderr io.Writer) *cobra.Command {
|
||||
func newCheckCommand(_ io.Reader, _, _ io.Writer) *cobra.Command {
|
||||
checker = ctl.NewCheckCommand(os.Stdin, os.Stdout, os.Stderr)
|
||||
checkCmd := &cobra.Command{
|
||||
Use: "check <path> [path2]...",
|
||||
|
|
@ -48,7 +48,3 @@ Performs a consistency check on data files.
|
|||
}
|
||||
return checkCmd
|
||||
}
|
||||
|
||||
func init() {
|
||||
subcommandFns["check"] = newCheckCommand
|
||||
}
|
||||
|
|
|
|||
|
|
@ -49,7 +49,3 @@ func newConfigCommand(stdin io.Reader, stdout, stderr io.Writer) *cobra.Command
|
|||
|
||||
return confCmd
|
||||
}
|
||||
|
||||
func init() {
|
||||
subcommandFns["config"] = newConfigCommand
|
||||
}
|
||||
|
|
|
|||
|
|
@ -26,7 +26,7 @@ import (
|
|||
|
||||
var Exporter *ctl.ExportCommand
|
||||
|
||||
func newExportCommand(stdin io.Reader, stdout, stderr io.Writer) *cobra.Command {
|
||||
func newExportCommand(_ io.Reader, _, _ io.Writer) *cobra.Command {
|
||||
Exporter = ctl.NewExportCommand(os.Stdin, os.Stdout, os.Stderr)
|
||||
exportCmd := &cobra.Command{
|
||||
Use: "export",
|
||||
|
|
@ -58,7 +58,3 @@ The file does not contain any headers.
|
|||
|
||||
return exportCmd
|
||||
}
|
||||
|
||||
func init() {
|
||||
subcommandFns["export"] = newExportCommand
|
||||
}
|
||||
|
|
|
|||
|
|
@ -26,7 +26,7 @@ import (
|
|||
|
||||
var generateConf *ctl.GenerateConfigCommand
|
||||
|
||||
func newGenerateConfigCommand(stdin io.Reader, stdout, stderr io.Writer) *cobra.Command {
|
||||
func newGenerateConfigCommand(_ io.Reader, _, _ io.Writer) *cobra.Command {
|
||||
generateConf = ctl.NewGenerateConfigCommand(os.Stdin, os.Stdout, os.Stderr)
|
||||
confCmd := &cobra.Command{
|
||||
Use: "generate-config",
|
||||
|
|
@ -43,7 +43,3 @@ func newGenerateConfigCommand(stdin io.Reader, stdout, stderr io.Writer) *cobra.
|
|||
|
||||
return confCmd
|
||||
}
|
||||
|
||||
func init() {
|
||||
subcommandFns["generate-config"] = newGenerateConfigCommand
|
||||
}
|
||||
|
|
|
|||
|
|
@ -65,7 +65,3 @@ omitted. If it is present then its format should be YYYY-MM-DDTHH:MM.
|
|||
|
||||
return importCmd
|
||||
}
|
||||
|
||||
func init() {
|
||||
subcommandFns["import"] = newImportCommand
|
||||
}
|
||||
|
|
|
|||
|
|
@ -27,7 +27,7 @@ import (
|
|||
|
||||
var inspector *ctl.InspectCommand
|
||||
|
||||
func newInspectCommand(stdin io.Reader, stdout, stderr io.Writer) *cobra.Command {
|
||||
func newInspectCommand(_ io.Reader, _, _ io.Writer) *cobra.Command {
|
||||
inspector = ctl.NewInspectCommand(os.Stdin, os.Stdout, os.Stderr)
|
||||
|
||||
inspectCmd := &cobra.Command{
|
||||
|
|
@ -51,7 +51,3 @@ Inspects a data file and provides stats.
|
|||
}
|
||||
return inspectCmd
|
||||
}
|
||||
|
||||
func init() {
|
||||
subcommandFns["inspect"] = newInspectCommand
|
||||
}
|
||||
|
|
|
|||
16
cmd/root.go
16
cmd/root.go
|
|
@ -25,10 +25,6 @@ import (
|
|||
"github.com/spf13/viper"
|
||||
)
|
||||
|
||||
// TODO maybe give this an Add method which will ensure two command
|
||||
// with same name aren't added
|
||||
var subcommandFns = map[string]func(stdin io.Reader, stdout, stderr io.Writer) *cobra.Command{}
|
||||
|
||||
func NewRootCommand(stdin io.Reader, stdout, stderr io.Writer) *cobra.Command {
|
||||
productName := "Pilosa " + pilosa.Version
|
||||
if pilosa.EnterpriseEnabled {
|
||||
|
|
@ -69,9 +65,15 @@ Build Time: ` + pilosa.BuildTime + "\n",
|
|||
rc.PersistentFlags().Bool("dry-run", false, "stop before executing")
|
||||
_ = rc.PersistentFlags().MarkHidden("dry-run")
|
||||
rc.PersistentFlags().StringP("config", "c", "", "Configuration file to read from.")
|
||||
for _, subcomFn := range subcommandFns {
|
||||
rc.AddCommand(subcomFn(stdin, stdout, stderr))
|
||||
}
|
||||
|
||||
rc.AddCommand(newCheckCommand(stdin, stdout, stderr))
|
||||
rc.AddCommand(newConfigCommand(stdin, stdout, stderr))
|
||||
rc.AddCommand(newExportCommand(stdin, stdout, stderr))
|
||||
rc.AddCommand(newGenerateConfigCommand(stdin, stdout, stderr))
|
||||
rc.AddCommand(newImportCommand(stdin, stdout, stderr))
|
||||
rc.AddCommand(newInspectCommand(stdin, stdout, stderr))
|
||||
rc.AddCommand(newServeCmd(stdin, stdout, stderr))
|
||||
|
||||
rc.SetOutput(stderr)
|
||||
return rc
|
||||
}
|
||||
|
|
|
|||
|
|
@ -50,7 +50,3 @@ on the configured port.`,
|
|||
ctl.BuildServerFlags(serveCmd, Server)
|
||||
return serveCmd
|
||||
}
|
||||
|
||||
func init() {
|
||||
subcommandFns["server"] = newServeCmd
|
||||
}
|
||||
|
|
|
|||
|
|
@ -203,6 +203,9 @@ func TestImportCommand_BugOverwriteValue(t *testing.T) {
|
|||
|
||||
file.Close()
|
||||
file, err = ioutil.TempFile("", "import-value2.csv")
|
||||
if err != nil {
|
||||
t.Fatalf("Error creating tempfile: %s", err)
|
||||
}
|
||||
file.Write([]byte("0,16\n"))
|
||||
cm.Paths = []string{file.Name()}
|
||||
err = cm.Run(ctx)
|
||||
|
|
@ -212,6 +215,9 @@ func TestImportCommand_BugOverwriteValue(t *testing.T) {
|
|||
|
||||
file.Close()
|
||||
file, err = ioutil.TempFile("", "import-value3.csv")
|
||||
if err != nil {
|
||||
t.Fatalf("Error creating tempfile: %s", err)
|
||||
}
|
||||
file.Write([]byte("0,19\n"))
|
||||
cm.Paths = []string{file.Name()}
|
||||
err = cm.Run(ctx)
|
||||
|
|
|
|||
|
|
@ -31,6 +31,9 @@ func TestInspectCommand_Run(t *testing.T) {
|
|||
|
||||
cm := NewInspectCommand(stdin, w, w)
|
||||
file, err := ioutil.TempFile("", "inspectTest")
|
||||
if err != nil {
|
||||
t.Fatalf("Error creating tempfile: %s", err)
|
||||
}
|
||||
file.Write([]byte("12358267538963"))
|
||||
file.Close()
|
||||
cm.Path = file.Name()
|
||||
|
|
|
|||
|
|
@ -32,7 +32,6 @@
|
|||
package b
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
"sync"
|
||||
|
||||
|
|
@ -40,20 +39,12 @@ import (
|
|||
)
|
||||
|
||||
const (
|
||||
// kx must be >= 2
|
||||
kx = 128 //TODO benchmark tune this number if using custom key/value type(s).
|
||||
// kd must be >= 1
|
||||
kd = 128 //TODO benchmark tune this number if using custom key/value type(s).
|
||||
)
|
||||
|
||||
func init() {
|
||||
if kd < 1 {
|
||||
panic(fmt.Errorf("kd %d: out of range", kd))
|
||||
}
|
||||
|
||||
if kx < 2 {
|
||||
panic(fmt.Errorf("kx %d: out of range", kx))
|
||||
}
|
||||
}
|
||||
|
||||
var (
|
||||
btDPool = sync.Pool{New: func() interface{} { return &d{} }}
|
||||
btEPool = btEpool{sync.Pool{New: func() interface{} { return &enumerator{} }}}
|
||||
|
|
@ -202,7 +193,7 @@ func (q *x) siblings(i int) (l, r *d) {
|
|||
r = q.x[i+1].ch.(*d)
|
||||
}
|
||||
}
|
||||
return
|
||||
return l, r
|
||||
}
|
||||
|
||||
// -------------------------------------------------------------------------- d
|
||||
|
|
@ -424,7 +415,7 @@ func (t *tree) First() (k uint64, v *roaring.Container) {
|
|||
q := &q.d[0]
|
||||
k, v = q.k, q.v
|
||||
}
|
||||
return
|
||||
return k, v
|
||||
}
|
||||
|
||||
// Get returns the value associated with k and true if it exists. Otherwise Get
|
||||
|
|
@ -475,7 +466,7 @@ func (t *tree) Last() (k uint64, v *roaring.Container) {
|
|||
q := &q.d[q.c-1]
|
||||
k, v = q.k, q.v
|
||||
}
|
||||
return
|
||||
return k, v
|
||||
}
|
||||
|
||||
// Len returns the number of items in the tree.
|
||||
|
|
@ -860,7 +851,7 @@ func (e *enumerator) Close() {
|
|||
// io.EOF is returned.
|
||||
func (e *enumerator) Next() (k uint64, v *roaring.Container, err error) {
|
||||
if err = e.err; err != nil {
|
||||
return
|
||||
return 0, nil, err
|
||||
}
|
||||
|
||||
if e.ver != e.t.ver {
|
||||
|
|
@ -870,12 +861,12 @@ func (e *enumerator) Next() (k uint64, v *roaring.Container, err error) {
|
|||
}
|
||||
if e.q == nil {
|
||||
e.err, err = io.EOF, io.EOF
|
||||
return
|
||||
return 0, nil, err
|
||||
}
|
||||
|
||||
if e.i >= e.q.c {
|
||||
if err = e.next(); err != nil {
|
||||
return
|
||||
return 0, nil, err
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -883,7 +874,7 @@ func (e *enumerator) Next() (k uint64, v *roaring.Container, err error) {
|
|||
k, v = i.k, i.v
|
||||
e.k, e.hit = k, true
|
||||
e.next()
|
||||
return
|
||||
return k, v, nil
|
||||
}
|
||||
|
||||
func (e *enumerator) next() error {
|
||||
|
|
@ -908,7 +899,7 @@ func (e *enumerator) next() error {
|
|||
// == io.EOF is returned.
|
||||
func (e *enumerator) Prev() (k uint64, v *roaring.Container, err error) {
|
||||
if err = e.err; err != nil {
|
||||
return
|
||||
return 0, nil, err
|
||||
}
|
||||
|
||||
if e.ver != e.t.ver {
|
||||
|
|
@ -918,19 +909,19 @@ func (e *enumerator) Prev() (k uint64, v *roaring.Container, err error) {
|
|||
}
|
||||
if e.q == nil {
|
||||
e.err, err = io.EOF, io.EOF
|
||||
return
|
||||
return 0, nil, err
|
||||
}
|
||||
|
||||
if !e.hit {
|
||||
// move to previous because Seek overshoots if there's no hit
|
||||
if err = e.prev(); err != nil {
|
||||
return
|
||||
return 0, nil, err
|
||||
}
|
||||
}
|
||||
|
||||
if e.i >= e.q.c {
|
||||
if err = e.prev(); err != nil {
|
||||
return
|
||||
return 0, nil, err
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -938,7 +929,7 @@ func (e *enumerator) Prev() (k uint64, v *roaring.Container, err error) {
|
|||
k, v = i.k, i.v
|
||||
e.k, e.hit = k, true
|
||||
e.prev()
|
||||
return
|
||||
return k, v, err
|
||||
}
|
||||
|
||||
func (e *enumerator) prev() error {
|
||||
|
|
|
|||
|
|
@ -128,7 +128,7 @@ func (btc *bTreeContainers) Count() (n uint64) {
|
|||
n += uint64(c.N())
|
||||
_, c, err = e.Next()
|
||||
}
|
||||
return
|
||||
return n
|
||||
}
|
||||
|
||||
func (btc *bTreeContainers) Clone() roaring.Containers {
|
||||
|
|
|
|||
|
|
@ -26,7 +26,7 @@ import (
|
|||
"github.com/pilosa/pilosa/roaring"
|
||||
)
|
||||
|
||||
func init() {
|
||||
func init() { // nolint: gochecknoinits
|
||||
// Replace Bitmap constructor with B+Tree implementation
|
||||
roaring.NewFileBitmap = b.NewBTreeBitmap
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1672,15 +1672,6 @@ type execOptions struct {
|
|||
ExcludeColumns bool
|
||||
}
|
||||
|
||||
// decodeError returns an error representation of s if s is non-blank.
|
||||
// Returns nil if s is blank.
|
||||
func decodeError(s string) error {
|
||||
if s == "" {
|
||||
return nil
|
||||
}
|
||||
return errors.New(s)
|
||||
}
|
||||
|
||||
// hasOnlySetRowAttrs returns true if calls only contains SetRowAttrs() calls.
|
||||
func hasOnlySetRowAttrs(calls []*pql.Call) bool {
|
||||
if len(calls) == 0 {
|
||||
|
|
|
|||
4
field.go
4
field.go
|
|
@ -821,9 +821,9 @@ func (f *Field) allTimeViewsSortedByQuantum() (me []*view) {
|
|||
}
|
||||
}
|
||||
}
|
||||
return
|
||||
return lt
|
||||
})
|
||||
return
|
||||
return me
|
||||
}
|
||||
|
||||
// Value reads a field value for a column.
|
||||
|
|
|
|||
80
fragment.go
80
fragment.go
|
|
@ -47,6 +47,13 @@ const (
|
|||
// ShardWidth is the number of column IDs in a shard.
|
||||
ShardWidth = 1048576
|
||||
|
||||
// containersPerRowSegment is dependent upon ShardWidth,
|
||||
// and it represents the number of containers per shard row
|
||||
// (or rowSegment). Since containers are set in roaring
|
||||
// to be 2^16, then this const should be ShardWidth / 2^16.
|
||||
// It is represented as the exponent n of 2^n.
|
||||
containersPerRowSegment = 4
|
||||
|
||||
// snapshotExt is the file extension used for an in-process snapshot.
|
||||
snapshotExt = ".snapshotting"
|
||||
|
||||
|
|
@ -153,7 +160,6 @@ func (f *fragment) Open() error {
|
|||
pos := f.storage.Max()
|
||||
f.maxRowID = pos / ShardWidth
|
||||
f.stats.Gauge("rows", float64(f.maxRowID), 1.0)
|
||||
|
||||
return nil
|
||||
}(); err != nil {
|
||||
f.close()
|
||||
|
|
@ -578,9 +584,9 @@ func (f *fragment) sum(filter *Row, bitDepth uint) (sum, count uint64, err error
|
|||
//
|
||||
// 10*(2^0) + 4*(2^1) + 3*(2^2) = 30
|
||||
//
|
||||
var cnt uint64
|
||||
for i := uint(0); i < bitDepth; i++ {
|
||||
row := f.row(uint64(i))
|
||||
cnt := uint64(0)
|
||||
if filter != nil {
|
||||
cnt = row.intersectionCount(filter)
|
||||
} else {
|
||||
|
|
@ -1159,12 +1165,11 @@ func (f *fragment) readContiguousChecksums(a *[]FragmentBlock, blockID int) (n i
|
|||
func (f *fragment) blockData(id int) (rowIDs, columnIDs []uint64) {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
f.storage.ForEachRange(uint64(id)*HashBlockSize*ShardWidth, (uint64(id)+1)*HashBlockSize*ShardWidth, func(i uint64) {
|
||||
rowIDs = append(rowIDs, i/ShardWidth)
|
||||
columnIDs = append(columnIDs, i%ShardWidth)
|
||||
})
|
||||
return
|
||||
return rowIDs, columnIDs
|
||||
}
|
||||
|
||||
// mergeBlock compares the block's bits and computes a diff with another set of block bits.
|
||||
|
|
@ -1680,6 +1685,63 @@ func (f *fragment) readCacheFromArchive(r io.Reader) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
func (f *fragment) rows() []uint64 {
|
||||
i, _ := f.storage.Containers.Iterator(0)
|
||||
rows := make([]uint64, 0)
|
||||
|
||||
var lastRow uint64
|
||||
lastRow = math.MaxUint64
|
||||
|
||||
// Loop over the existing containers.
|
||||
for i.Next() {
|
||||
key, _ := i.Value()
|
||||
|
||||
// virtual row for the current container
|
||||
vRow := key >> containersPerRowSegment
|
||||
|
||||
// skip dups
|
||||
if vRow == lastRow {
|
||||
continue
|
||||
}
|
||||
|
||||
rows = append(rows, vRow)
|
||||
lastRow = vRow
|
||||
}
|
||||
return rows
|
||||
|
||||
}
|
||||
|
||||
func (f *fragment) rowsForColumn(columnID uint64) []uint64 {
|
||||
var colKey uint64
|
||||
|
||||
colID := columnID % ShardWidth
|
||||
i, _ := f.storage.Containers.Iterator(0)
|
||||
|
||||
colVal := uint16(colID & 0xFFFF)
|
||||
|
||||
rows := make([]uint64, 0)
|
||||
|
||||
// Loop over the existing containers.
|
||||
for i.Next() {
|
||||
key, c := i.Value()
|
||||
|
||||
// virtual row for the current container
|
||||
vRow := key >> containersPerRowSegment
|
||||
|
||||
// column container key for virtual row
|
||||
colKey = ((vRow * ShardWidth) + colID) >> 16
|
||||
|
||||
if colKey != key {
|
||||
continue
|
||||
}
|
||||
|
||||
if c.Contains(colVal) {
|
||||
rows = append(rows, vRow)
|
||||
}
|
||||
}
|
||||
return rows
|
||||
}
|
||||
|
||||
// FragmentBlock represents info about a subsection of the rows in a block.
|
||||
// This is used for comparing data in remote blocks for active anti-entropy.
|
||||
type FragmentBlock struct {
|
||||
|
|
@ -1903,12 +1965,12 @@ func (s *fragmentSyncer) syncBlock(id int) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
func madvise(b []byte, advice int) (err error) { // nolint: unparam
|
||||
_, _, e1 := syscall.Syscall(syscall.SYS_MADVISE, uintptr(unsafe.Pointer(&b[0])), uintptr(len(b)), uintptr(advice))
|
||||
if e1 != 0 {
|
||||
err = e1
|
||||
func madvise(b []byte, advice int) error { // nolint: unparam
|
||||
_, _, err := syscall.Syscall(syscall.SYS_MADVISE, uintptr(unsafe.Pointer(&b[0])), uintptr(len(b)), uintptr(advice))
|
||||
if err != 0 {
|
||||
return err
|
||||
}
|
||||
return
|
||||
return nil
|
||||
}
|
||||
|
||||
// pairSet is a list of equal length row and column id lists.
|
||||
|
|
|
|||
|
|
@ -1250,7 +1250,9 @@ func mustOpenFragment(index, field, view string, shard uint64, cacheType string)
|
|||
|
||||
f := newFragment(file.Name(), index, field, view, shard)
|
||||
f.CacheType = cacheType
|
||||
f.RowAttrStore = newMemAttrStore()
|
||||
f.RowAttrStore = &memAttrStore{
|
||||
store: make(map[uint64]map[string]interface{}),
|
||||
}
|
||||
|
||||
if err := f.Open(); err != nil {
|
||||
panic(err)
|
||||
|
|
@ -1278,3 +1280,81 @@ func (f *fragment) mustSetBits(rowID uint64, columnIDs ...uint64) {
|
|||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Test Various methods of retrieving RowIDs
|
||||
func TestFragment_RowsIteration(t *testing.T) {
|
||||
t.Run("firstContainer", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
|
||||
expectedAll := make([]uint64, 0)
|
||||
expectedOdd := make([]uint64, 0)
|
||||
for i := uint64(100); i < uint64(200); i++ {
|
||||
if _, err := f.setBit(i, i%2); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
expectedAll = append(expectedAll, i)
|
||||
if i%2 == 1 {
|
||||
expectedOdd = append(expectedOdd, i)
|
||||
}
|
||||
}
|
||||
|
||||
ids := f.rows()
|
||||
if !reflect.DeepEqual(expectedAll, ids) {
|
||||
t.Fatalf("Do not match %v %v", expectedAll, ids)
|
||||
}
|
||||
|
||||
ids = f.rowsForColumn(1)
|
||||
if !reflect.DeepEqual(expectedOdd, ids) {
|
||||
t.Fatalf("Do not match %v %v", expectedOdd, ids)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("secondRow", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
|
||||
expected := []uint64{1, 2}
|
||||
if _, err := f.setBit(1, 66000); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.setBit(2, 66000); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.setBit(2, 166000); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
ids := f.rows()
|
||||
if !reflect.DeepEqual(expected, ids) {
|
||||
t.Fatalf("Do not match %v %v", expected, ids)
|
||||
}
|
||||
|
||||
ids = f.rowsForColumn(66000)
|
||||
if !reflect.DeepEqual(expected, ids) {
|
||||
t.Fatalf("Do not match %v %v", expected, ids)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("combinations", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
|
||||
expectedRows := make([]uint64, 0)
|
||||
for r := uint64(1); r < uint64(10000); r += 100 {
|
||||
expectedRows = append(expectedRows, r)
|
||||
for c := uint64(1); c < uint64(ShardWidth-1); c += 10000 {
|
||||
if _, err := f.setBit(r, c); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
ids := f.rows()
|
||||
if !reflect.DeepEqual(expectedRows, ids) {
|
||||
t.Fatalf("Do not match %v %v", expectedRows, ids)
|
||||
}
|
||||
ids = f.rowsForColumn(c)
|
||||
if !reflect.DeepEqual(expectedRows, ids) {
|
||||
t.Fatalf("Do not match %v %v", expectedRows, ids)
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
|
|
|||
6
gc.go
6
gc.go
|
|
@ -23,12 +23,8 @@ type GCNotifier interface {
|
|||
AfterGC() <-chan struct{}
|
||||
}
|
||||
|
||||
func init() {
|
||||
NopGCNotifier = &nopGCNotifier{}
|
||||
}
|
||||
|
||||
// NopGCNotifier represents a GCNotifier that doesn't do anything.
|
||||
var NopGCNotifier GCNotifier
|
||||
var NopGCNotifier GCNotifier = &nopGCNotifier{}
|
||||
|
||||
type nopGCNotifier struct{}
|
||||
|
||||
|
|
|
|||
|
|
@ -29,13 +29,6 @@ import (
|
|||
"github.com/pilosa/pilosa/test"
|
||||
)
|
||||
|
||||
var defaultClient *gohttp.Client
|
||||
|
||||
func init() {
|
||||
defaultClient = http.GetHTTPClient(nil)
|
||||
|
||||
}
|
||||
|
||||
// Test distributed TopN Row count across 3 nodes.
|
||||
func TestClient_MultiNode(t *testing.T) {
|
||||
c := test.MustRunCluster(t, 3,
|
||||
|
|
@ -125,9 +118,9 @@ func TestClient_MultiNode(t *testing.T) {
|
|||
|
||||
// Connect to each node to compare results.
|
||||
client := make([]*Client, 3)
|
||||
client[0] = MustNewClient(c[0].URL(), defaultClient)
|
||||
client[1] = MustNewClient(c[1].URL(), defaultClient)
|
||||
client[2] = MustNewClient(c[2].URL(), defaultClient)
|
||||
client[0] = MustNewClient(c[0].URL(), http.GetHTTPClient(nil))
|
||||
client[1] = MustNewClient(c[1].URL(), http.GetHTTPClient(nil))
|
||||
client[2] = MustNewClient(c[2].URL(), http.GetHTTPClient(nil))
|
||||
|
||||
topN := 4
|
||||
queryRequest := &pilosa.QueryRequest{
|
||||
|
|
@ -191,7 +184,7 @@ func TestClient_Import(t *testing.T) {
|
|||
hldr.Row("i", "f", 0)
|
||||
|
||||
// Send import request.
|
||||
c := MustNewClient(host, defaultClient)
|
||||
c := MustNewClient(host, http.GetHTTPClient(nil))
|
||||
if err := c.Import(context.Background(), "i", "f", 0, []pilosa.Bit{
|
||||
{RowID: 0, ColumnID: 1},
|
||||
{RowID: 0, ColumnID: 5},
|
||||
|
|
@ -226,7 +219,7 @@ func TestClient_ImportValue(t *testing.T) {
|
|||
}
|
||||
|
||||
// Send import request.
|
||||
c := MustNewClient(host, defaultClient)
|
||||
c := MustNewClient(host, http.GetHTTPClient(nil))
|
||||
if err := c.ImportValue(context.Background(), "i", "f", 0, []pilosa.FieldValue{
|
||||
{ColumnID: 1, Value: -10},
|
||||
{ColumnID: 2, Value: 20},
|
||||
|
|
@ -287,7 +280,7 @@ func TestClient_FragmentBlocks(t *testing.T) {
|
|||
|
||||
// Set a bit on a different shard.
|
||||
hldr.SetBit("i", "f", 0, 1)
|
||||
c := MustNewClient(cmd.URL(), defaultClient)
|
||||
c := MustNewClient(cmd.URL(), http.GetHTTPClient(nil))
|
||||
blocks, err := c.FragmentBlocks(context.Background(), nil, "i", "f", 0)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
|
|
|
|||
|
|
@ -302,7 +302,7 @@ type successResponse struct {
|
|||
func (r *successResponse) check(err error) (statusCode int) {
|
||||
if err == nil {
|
||||
r.Success = true
|
||||
return
|
||||
return 0
|
||||
}
|
||||
|
||||
cause := errors.Cause(err)
|
||||
|
|
@ -322,7 +322,7 @@ func (r *successResponse) check(err error) (statusCode int) {
|
|||
r.Success = false
|
||||
r.Error = &Error{Message: cause.Error()}
|
||||
|
||||
return
|
||||
return statusCode
|
||||
}
|
||||
|
||||
// write sends a response to the http.ResponseWriter based on the success
|
||||
|
|
@ -1122,12 +1122,9 @@ func (h *Handler) handleGetVersion(w http.ResponseWriter, r *http.Request) {
|
|||
|
||||
// QueryResult types.
|
||||
const (
|
||||
queryResultTypeNil uint32 = iota
|
||||
QueryResultTypeRow
|
||||
QueryResultTypeRow uint32 = iota
|
||||
QueryResultTypePairs
|
||||
queryResultTypeValCount
|
||||
QueryResultTypeUint64
|
||||
queryResultTypeBool
|
||||
)
|
||||
|
||||
// parseUint64Slice returns a slice of uint64s from a comma-delimited string.
|
||||
|
|
@ -1149,14 +1146,6 @@ func parseUint64Slice(s string) ([]uint64, error) {
|
|||
return a, nil
|
||||
}
|
||||
|
||||
// errorString returns the string representation of err.
|
||||
func errorString(err error) string {
|
||||
if err == nil {
|
||||
return ""
|
||||
}
|
||||
return err.Error()
|
||||
}
|
||||
|
||||
func (h *Handler) handlePostClusterResizeSetCoordinator(w http.ResponseWriter, r *http.Request) {
|
||||
if !validHeaderAcceptJSON(r.Header) {
|
||||
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
|
||||
|
|
|
|||
|
|
@ -28,12 +28,8 @@ type Logger interface {
|
|||
Debugf(format string, v ...interface{})
|
||||
}
|
||||
|
||||
func init() {
|
||||
NopLogger = &nopLogger{}
|
||||
}
|
||||
|
||||
// NopLogger represents a Logger that doesn't do anything.
|
||||
var NopLogger Logger
|
||||
var NopLogger Logger = &nopLogger{}
|
||||
|
||||
type nopLogger struct{}
|
||||
|
||||
|
|
|
|||
|
|
@ -73,13 +73,13 @@ func (c *Cache) Add(key Key, value interface{}) {
|
|||
// Get looks up a key's value from the cache.
|
||||
func (c *Cache) Get(key Key) (value interface{}, ok bool) {
|
||||
if c.cache == nil {
|
||||
return
|
||||
return nil, false
|
||||
}
|
||||
if ele, hit := c.cache[key]; hit {
|
||||
c.ll.MoveToFront(ele)
|
||||
return ele.Value.(*entry).value, true
|
||||
}
|
||||
return
|
||||
return nil, false
|
||||
}
|
||||
|
||||
// remove removes the provided key from the cache.
|
||||
|
|
|
|||
|
|
@ -177,7 +177,7 @@ func (b *Bitmap) Contains(v uint64) bool {
|
|||
if c == nil {
|
||||
return false
|
||||
}
|
||||
return c.contains(lowbits(v))
|
||||
return c.Contains(lowbits(v))
|
||||
}
|
||||
|
||||
// Remove removes values from the bitmap.
|
||||
|
|
@ -1272,8 +1272,8 @@ func (c *Container) runAdd(v uint16) bool {
|
|||
return true
|
||||
}
|
||||
|
||||
// contains returns true if v is in the container.
|
||||
func (c *Container) contains(v uint16) bool {
|
||||
// Contains returns true if v is in the container.
|
||||
func (c *Container) Contains(v uint16) bool {
|
||||
if c.isArray() {
|
||||
return c.arrayContains(v)
|
||||
} else if c.isRun() {
|
||||
|
|
@ -1800,7 +1800,7 @@ type containerInfo struct {
|
|||
|
||||
// flip returns a new container containing the inverse of all
|
||||
// bits in a.
|
||||
func flip(a *Container) *Container {
|
||||
func flip(a *Container) *Container { // nolint: deadcode
|
||||
if a.isArray() {
|
||||
return flipArray(a)
|
||||
} else if a.isRun() {
|
||||
|
|
@ -1923,7 +1923,7 @@ func intersectionCountRunRun(a, b *Container) (n int) {
|
|||
i++
|
||||
}
|
||||
}
|
||||
return
|
||||
return n
|
||||
}
|
||||
|
||||
func intersectionCountBitmapRun(a, b *Container) (n int) {
|
||||
|
|
@ -3149,17 +3149,13 @@ func xorCompare(x *xorstm) (r1 interval16, hasData bool) {
|
|||
if !x.vaValid || !x.vbValid {
|
||||
if x.vbValid {
|
||||
x.vbValid = false
|
||||
r1 = x.vb
|
||||
hasData = true
|
||||
return
|
||||
return x.vb, true
|
||||
}
|
||||
if x.vaValid {
|
||||
x.vaValid = false
|
||||
r1 = x.va
|
||||
hasData = true
|
||||
return
|
||||
return x.va, true
|
||||
}
|
||||
return
|
||||
return r1, false
|
||||
}
|
||||
|
||||
if x.va.last < x.vb.start { //va before
|
||||
|
|
@ -3232,7 +3228,7 @@ func xorCompare(x *xorstm) (r1 interval16, hasData bool) {
|
|||
}
|
||||
}
|
||||
}
|
||||
return
|
||||
return r1, hasData
|
||||
}
|
||||
|
||||
//stm is state machine used to "xor" iterate over runs.
|
||||
|
|
@ -3304,7 +3300,7 @@ func xorBitmapRun(a, b *Container) *Container {
|
|||
return output
|
||||
}
|
||||
|
||||
func bitmapsEqual(b, c *Bitmap) error {
|
||||
func bitmapsEqual(b, c *Bitmap) error { // nolint: deadcode
|
||||
if b.OpWriter != c.OpWriter {
|
||||
return errors.New("opWriters not equal")
|
||||
}
|
||||
|
|
@ -3336,22 +3332,6 @@ func popcount(x uint64) uint64 {
|
|||
return uint64(bits.OnesCount64(x))
|
||||
}
|
||||
|
||||
func popcountSlice(s []uint64) uint64 {
|
||||
cnt := uint64(0)
|
||||
for _, x := range s {
|
||||
cnt += popcount(x)
|
||||
}
|
||||
return cnt
|
||||
}
|
||||
|
||||
func popcountMaskSlice(s, m []uint64) uint64 {
|
||||
cnt := uint64(0)
|
||||
for i := range s {
|
||||
cnt += popcount(s[i] &^ m[i])
|
||||
}
|
||||
return cnt
|
||||
}
|
||||
|
||||
func popcountAndSlice(s, m []uint64) uint64 {
|
||||
cnt := uint64(0)
|
||||
for i := range s {
|
||||
|
|
@ -3359,19 +3339,3 @@ func popcountAndSlice(s, m []uint64) uint64 {
|
|||
}
|
||||
return cnt
|
||||
}
|
||||
|
||||
func popcountOrSlice(s, m []uint64) uint64 {
|
||||
cnt := uint64(0)
|
||||
for i := range s {
|
||||
cnt += popcount(s[i] | m[i])
|
||||
}
|
||||
return cnt
|
||||
}
|
||||
|
||||
func popcountXorSlice(s, m []uint64) uint64 {
|
||||
cnt := uint64(0)
|
||||
for i := range s {
|
||||
cnt += popcount(s[i] ^ m[i])
|
||||
}
|
||||
return cnt
|
||||
}
|
||||
|
|
|
|||
|
|
@ -2276,7 +2276,7 @@ func TestIteratorRuns(t *testing.T) {
|
|||
t.Fatalf("iterator did not seek correctly in multiple containers: %v\n", itr)
|
||||
}
|
||||
|
||||
val, eof = itr.Next()
|
||||
itr.Next()
|
||||
val, eof = itr.Next()
|
||||
if !(val == 0 && eof) {
|
||||
t.Fatalf("iterator did not eof correctly: %d, %v\n", val, eof)
|
||||
|
|
|
|||
|
|
@ -44,10 +44,6 @@ import (
|
|||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
func init() {
|
||||
rand.Seed(time.Now().UTC().UnixNano())
|
||||
}
|
||||
|
||||
type loggerLogger interface {
|
||||
pilosa.Logger
|
||||
Logger() *log.Logger
|
||||
|
|
@ -126,6 +122,9 @@ func NewCommand(stdin io.Reader, stdout, stderr io.Writer, opts ...CommandOption
|
|||
func (m *Command) Start() (err error) {
|
||||
defer close(m.Started)
|
||||
|
||||
// Seed random number generator
|
||||
rand.Seed(time.Now().UTC().UnixNano())
|
||||
|
||||
// SetupServer
|
||||
err = m.SetupServer()
|
||||
if err != nil {
|
||||
|
|
|
|||
6
stats.go
6
stats.go
|
|
@ -22,10 +22,6 @@ import (
|
|||
"time"
|
||||
)
|
||||
|
||||
func init() {
|
||||
NopStatsClient = &nopStatsClient{}
|
||||
}
|
||||
|
||||
// Expvar global expvar map.
|
||||
var Expvar = expvar.NewMap("index")
|
||||
|
||||
|
|
@ -66,7 +62,7 @@ type StatsClient interface {
|
|||
}
|
||||
|
||||
// NopStatsClient represents a client that doesn't do anything.
|
||||
var NopStatsClient StatsClient
|
||||
var NopStatsClient StatsClient = &nopStatsClient{}
|
||||
|
||||
type nopStatsClient struct{}
|
||||
|
||||
|
|
|
|||
|
|
@ -243,6 +243,9 @@ func MustDo(method, urlStr string, body string) *httpResponse {
|
|||
urlStr,
|
||||
strings.NewReader(body),
|
||||
)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("Accept", "application/json")
|
||||
|
|
|
|||
|
|
@ -40,6 +40,9 @@ func TestNewCluster(t *testing.T) {
|
|||
cluster[0].URL()+"/status",
|
||||
strings.NewReader(""),
|
||||
)
|
||||
if err != nil {
|
||||
t.Fatalf("creating http request: %v", err)
|
||||
}
|
||||
|
||||
req.Header.Set("Accept", "application/json")
|
||||
|
||||
|
|
|
|||
|
|
@ -19,7 +19,9 @@ var EnterpriseEnabled = false
|
|||
var Version = "v0.0.0"
|
||||
var BuildTime = "not recorded"
|
||||
|
||||
func init() {
|
||||
// init sets the EnterpriseEnabled bool, based on the Enterprise string.
|
||||
// This is needed because bools cannot be set with ldflags.
|
||||
func init() { // nolint: gochecknoinits
|
||||
if Enterprise == "1" {
|
||||
EnterpriseEnabled = true
|
||||
}
|
||||
|
|
|
|||
|
|
@ -30,7 +30,9 @@ func mustOpenView(index, field, name string) *view {
|
|||
if err := v.open(); err != nil {
|
||||
panic(err)
|
||||
}
|
||||
v.rowAttrStore = newMemAttrStore()
|
||||
v.rowAttrStore = &memAttrStore{
|
||||
store: make(map[uint64]map[string]interface{}),
|
||||
}
|
||||
return v
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue