mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-07 17:15:56 +00:00
* use t.Fatal(f) to abort tests, not panic * make perf_able run at all, make it debug a bit better switch perf-able to using same node type we use for other spot instances, because otherwise it never finds any available capacity. we switch the perf-able script to use the standard get_value function instead of direct jq calls. we try to grab server logs if the restore fails in the hopes of finding out why the restore very occasionally fails. * Fix some issues with running IDK tests in docker. (#2248) *Stop running TestKafkaSourceIntegration with t.Parallel() This test can't be run in parallel as it's currently written. Doing so allows for interleaving of messages to the same kafka topic between tests. I didn't attempt to modify the test so it could be run in parallel. That could be done, but left for someone more ambitious. * Remove idk/testenv/certs which got accidentally committed. also update .gitignore to include those. * changes to add bool support in idk (#2240) * initial changes to add bool support in idk * modifying some default parameters for testing, will revert them later * adding support for bool in making fragments function * boolean values implementation without supporting empty or null values at this point * Implement bool support in batch using a map (and a slice for nulls) (#2247) * Implement bool support in batch using a map (and a slice for nulls) * Keep the PackBools default for now But set it explicity in the ingest tests which rely on it. * Modify batch to construct bool update like mutex The code in API.ImportRoaringShard has a switch statement which causes bool fields to be handled like mutex fields. This means, that the viewUpdate.Clear value should only contain data in the first "row" of the fragment, which it will treat as records to clear for *all* rows. This makes more sense for mutex fields; for bool fields, there's only one other row to clear. But since the code is currently handling them the same, we need to construct viewUpdate.Clear such that it conforms to that pattern. This commit also adds a test which covers this logic. * Remove commented code; revert config for testing This commit also removes the DELETE_SENTINEL case for non-packed bools, since that isn't supported anyway. * Revert default setting * remove inconsistent type scope * correcting the logic of string converstion to bool * resolving an error in a test * adding tests to cover code related to bool support in batch.go file and interface.go files * modifying interfaces test * added one more test case Co-authored-by: Travis Turner <travis@pilosa.com> Co-authored-by: Travis Turner <travis@molecula.com> * resolving bool null field ingestion error (#2254) * resolving bool null field ingestion error * testing issues * adding null support for bools * updating the null bool field ingestion * trying to resolve issue when ingesting null value for bool type * adding a clearing support for bool type * resolving issues with bool null value ingestion * updating the jwt go package version and removing changes made in docker compose file * reverting jwt go version * removing v4 of jwt * adding a comment in test file to see if sonar cloud accepts this file * don't obtain stack traces on rbf.Tx creation We thought stack traces were mildly expensive. We were very wrong. Due to a complicated issue in the Go runtime, simultaneous requests for stack traces end up contending on a lock even when they're not actually contending on any resources. I've filed a ticket in the Go issue tracker for this: https://github.com/golang/go/issues/56400 In the mean time: Under some workloads, we were seeing 85% of all CPU time go into the stack backtraces, of which 81% went into the contention on those locks. But even if you take away the contention, that leaves us with 4/19 of all CPU time in our code going into building those stack backtraces. That's a lot of overhead for a feature we virtually never use. We might consider adding a backtrace functionality here, possibly using `runtime.Callers` which is much lower overhead, and allows us to generate a backtrace on demand (no argument values available, but then, we never read those because they're unformatted hex values), but I don't think it's actually very informative to know what the stack traces were of the Tx; they don't necessarily reflect the current state of any ongoing use of the Tx, so we can't necessarily correlate them to goroutine stack dumps, and so on. * fb-1729 Enriched Table Metadata (#2255) enriched metadata for tables added support for the concept of a table and field owners in metadata; mechanism to derive owner from http request metadata; metadata for table description * tightened up is/is not null filter expressions (FB-1741) (#2260) Covers tightening up handling filter expressions that contain is/is not null ops. These filters may have to be translated into PQL calls to be passed to the executor and even though sql3 language supports nullability for any data type, currently only BSI fields are nullable at the storage engine level (there is a ticket to add support for non-BSI field here FB-1689: IS SQL Argument returns incorrect error) so when these fields are used in filter conditions we need to handle BSI and non-BSI fields differently. * added a test to cover the keyword replace as being synonymous with insert (#2261) * update molecula references to featurebase (#2262) Co-authored-by: Seebs <seebs@molecula.com> Co-authored-by: Travis Turner <travis@pilosa.com> Co-authored-by: Pranitha-malae <56414132+Pranitha-malae@users.noreply.github.com> Co-authored-by: Travis Turner <travis@molecula.com> Co-authored-by: pokeeffe-molecula <85502298+pokeeffe-molecula@users.noreply.github.com> Co-authored-by: Stephanie Yang <stephanie@pilosa.com>
433 lines
13 KiB
Go
433 lines
13 KiB
Go
// Copyright 2022 Molecula Corp. (DBA FeatureBase).
|
|
// SPDX-License-Identifier: Apache-2.0
|
|
package disco
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"path"
|
|
"sync"
|
|
)
|
|
|
|
var (
|
|
ErrTooManyResults error = fmt.Errorf("too many results")
|
|
ErrNoResults error = fmt.Errorf("no results")
|
|
ErrKeyDeleted error = fmt.Errorf("key deleted")
|
|
ErrIndexExists error = fmt.Errorf("index already exists")
|
|
ErrIndexDoesNotExist error = fmt.Errorf("index does not exist")
|
|
ErrFieldExists error = fmt.Errorf("field already exists")
|
|
ErrFieldDoesNotExist error = fmt.Errorf("field does not exist")
|
|
ErrViewExists error = fmt.Errorf("view already exists")
|
|
ErrViewDoesNotExist error = fmt.Errorf("view does not exist")
|
|
ErrKeyDoesNotExist error = fmt.Errorf("key does not exist")
|
|
)
|
|
|
|
type Peer struct {
|
|
URL string
|
|
ID string
|
|
}
|
|
|
|
func (p *Peer) String() string {
|
|
return fmt.Sprintf(`{"ID": "%s", "URL": "%s"}`, p.ID, p.URL)
|
|
}
|
|
|
|
type DisCo interface {
|
|
io.Closer
|
|
|
|
Start(ctx context.Context) (InitialClusterState, error)
|
|
IsLeader() bool
|
|
ID() string
|
|
Leader() *Peer
|
|
Peers() []*Peer
|
|
DeleteNode(ctx context.Context, id string) error
|
|
}
|
|
|
|
type (
|
|
InitialClusterState string
|
|
|
|
// ClusterState represents the state returned in the /status endpoint.
|
|
ClusterState string
|
|
)
|
|
|
|
const (
|
|
InitialClusterStateNew InitialClusterState = "new"
|
|
InitialClusterStateExisting InitialClusterState = "existing"
|
|
|
|
ClusterStateUnknown ClusterState = "UNKNOWN" // default cluster state. It is returned when we are not able to get the real actual state.
|
|
ClusterStateStarting ClusterState = "STARTING" // cluster is starting and some internal services are not ready yet.
|
|
ClusterStateDegraded ClusterState = "DEGRADED" // cluster is running but we've lost some # of hosts >0 but < replicaN. Only read queries are allowed.
|
|
ClusterStateNormal ClusterState = "NORMAL" // cluster is up and running.
|
|
ClusterStateDown ClusterState = "DOWN" // cluster is unable to serve queries.
|
|
)
|
|
|
|
type NodeState string
|
|
|
|
const (
|
|
NodeStateUnknown NodeState = "UNKNOWN"
|
|
NodeStateStarting NodeState = "STARTING"
|
|
NodeStateStarted NodeState = "STARTED"
|
|
)
|
|
|
|
// Schema is a map of all indexes, each of those being a map of fields, then
|
|
// views.
|
|
type Schema map[string]*Index
|
|
|
|
// Index is a struct which contains the data encoded for the index as well as
|
|
// for each of its fields.
|
|
type Index struct {
|
|
Data []byte
|
|
Fields map[string]*Field
|
|
}
|
|
|
|
// Field is a struct which contains the data encoded for the field as well as
|
|
// for each of its views.
|
|
type Field struct {
|
|
Data []byte
|
|
Views map[string]struct{}
|
|
}
|
|
|
|
// Schemator is the source of truth for different schema elements.
|
|
// All nodes will store and retrieve information from the same source,
|
|
// having the same information at the same time.
|
|
type Schemator interface {
|
|
|
|
// Schema return the actual pilosa schema. If the schema is not present, an error is returned.
|
|
Schema(ctx context.Context) (Schema, error)
|
|
|
|
// Index gets a specific index data by name.
|
|
Index(ctx context.Context, name string) ([]byte, error)
|
|
|
|
CreateIndex(ctx context.Context, name string, val []byte) error
|
|
DeleteIndex(ctx context.Context, name string) error
|
|
Field(ctx context.Context, index, field string) ([]byte, error)
|
|
CreateField(ctx context.Context, index, field string, fieldVal []byte) error
|
|
UpdateField(ctx context.Context, index, field string, fieldVal []byte) error
|
|
DeleteField(ctx context.Context, index, field string) error
|
|
View(ctx context.Context, index, field, view string) (bool, error)
|
|
CreateView(ctx context.Context, index, field, view string) error
|
|
DeleteView(ctx context.Context, index, field, view string) error
|
|
}
|
|
|
|
// Sharder is an interface used to maintain the set of availableShards bitmaps
|
|
// per field.
|
|
type Sharder interface {
|
|
Shards(ctx context.Context, index, field string) ([][]byte, error)
|
|
SetShards(ctx context.Context, index, field string, shards []byte) error
|
|
}
|
|
|
|
// NopDisCo represents a DisCo that doesn't do anything.
|
|
var NopDisCo DisCo = &nopDisCo{}
|
|
|
|
type nopDisCo struct{}
|
|
|
|
// Close no-op.
|
|
func (n *nopDisCo) Close() error {
|
|
return nil
|
|
}
|
|
|
|
// Start is a no-op implementation of the DisCo Start method.
|
|
func (n *nopDisCo) Start(ctx context.Context) (InitialClusterState, error) {
|
|
return InitialClusterStateNew, nil
|
|
}
|
|
|
|
// ID is a no-op implementation of the DisCo ID method.
|
|
func (n *nopDisCo) ID() string {
|
|
return ""
|
|
}
|
|
|
|
// IsLeader is a no-op implementation of the DisCo IsLeader method.
|
|
func (n *nopDisCo) IsLeader() bool {
|
|
return false
|
|
}
|
|
|
|
// Leader is a no-op implementation of the DisCo Leader method.
|
|
func (n *nopDisCo) Leader() *Peer {
|
|
return nil
|
|
}
|
|
|
|
// Peers is a no-op implementation of the DisCo Peers method.
|
|
func (n *nopDisCo) Peers() []*Peer {
|
|
return nil
|
|
}
|
|
|
|
// DeleteNode a no-op implementation of the DisCo DeleteNode method.
|
|
func (n *nopDisCo) DeleteNode(context.Context, string) error {
|
|
return nil
|
|
}
|
|
|
|
// NopSharder represents a Sharder that doesn't do anything.
|
|
var NopSharder Sharder = &nopSharder{}
|
|
|
|
type nopSharder struct{}
|
|
|
|
// Shards is a no-op implementation of the Sharder Shards method.
|
|
func (n *nopSharder) Shards(ctx context.Context, index, field string) ([][]byte, error) {
|
|
return nil, nil
|
|
}
|
|
|
|
// AddShards is a no-op implementation of the Sharder AddShards method.
|
|
func (n *nopSharder) SetShards(ctx context.Context, index, field string, shards []byte) error {
|
|
return nil
|
|
}
|
|
|
|
// NopSchemator represents a Schemator that doesn't do anything.
|
|
var NopSchemator Schemator = &nopSchemator{}
|
|
|
|
type nopSchemator struct{}
|
|
|
|
// Schema is a no-op implementation of the Schemator Schema method.
|
|
func (*nopSchemator) Schema(ctx context.Context) (Schema, error) { return nil, nil }
|
|
|
|
// Index is a no-op implementation of the Schemator Index method.
|
|
func (*nopSchemator) Index(ctx context.Context, name string) ([]byte, error) { return nil, nil }
|
|
|
|
// CreateIndex is a no-op implementation of the Schemator CreateIndex method.
|
|
func (*nopSchemator) CreateIndex(ctx context.Context, name string, val []byte) error { return nil }
|
|
|
|
// DeleteIndex is a no-op implementation of the Schemator DeleteIndex method.
|
|
func (*nopSchemator) DeleteIndex(ctx context.Context, name string) error {
|
|
return nil
|
|
}
|
|
|
|
// Field is a no-op implementation of the Schemator Field method.
|
|
func (*nopSchemator) Field(ctx context.Context, index, field string) ([]byte, error) { return nil, nil }
|
|
|
|
// CreateField is a no-op implementation of the Schemator CreateField method.
|
|
func (*nopSchemator) CreateField(ctx context.Context, index, field string, fieldVal []byte) error {
|
|
return nil
|
|
}
|
|
|
|
// UpdateField is a no-op implementation of the Schemator UpdateField method.
|
|
func (*nopSchemator) UpdateField(ctx context.Context, index, field string, fieldVal []byte) error {
|
|
return nil
|
|
}
|
|
|
|
// DeleteField is a no-op implementation of the Schemator DeleteField method.
|
|
func (*nopSchemator) DeleteField(ctx context.Context, index, field string) error {
|
|
return nil
|
|
}
|
|
|
|
// View is a no-op implementation of the Schemator View method.
|
|
func (*nopSchemator) View(ctx context.Context, index, field, view string) (bool, error) {
|
|
return false, nil
|
|
}
|
|
|
|
// CreateView is a no-op implementation of the Schemator CreateView method.
|
|
func (*nopSchemator) CreateView(ctx context.Context, index, field, view string) error {
|
|
return nil
|
|
}
|
|
|
|
// DeleteView is a no-op implementation of the Schemator DeleteView method.
|
|
func (*nopSchemator) DeleteView(ctx context.Context, index, field, view string) error { return nil }
|
|
|
|
// InMemSchemator represents a Schemator that manages the schema in memory. The
|
|
// intention is that this would be used for testing.
|
|
var InMemSchemator Schemator = &inMemSchemator{
|
|
schema: make(Schema),
|
|
}
|
|
|
|
type inMemSchemator struct {
|
|
mu sync.RWMutex
|
|
schema Schema
|
|
}
|
|
|
|
// NewInMemSchemator instantiates an InMemSchemator
|
|
// this allows new holders to have thier own, and not rely on a shared instance
|
|
func NewInMemSchemator() *inMemSchemator {
|
|
return &inMemSchemator{
|
|
schema: make(Schema),
|
|
}
|
|
}
|
|
|
|
// Schema is an in-memory implementation of the Schemator Schema method.
|
|
func (s *inMemSchemator) Schema(ctx context.Context) (Schema, error) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return s.schema, nil
|
|
}
|
|
|
|
// Index is an in-memory implementation of the Schemator Index method.
|
|
func (s *inMemSchemator) Index(ctx context.Context, name string) ([]byte, error) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
idx, ok := s.schema[name]
|
|
if !ok {
|
|
return nil, ErrIndexDoesNotExist
|
|
}
|
|
return idx.Data, nil
|
|
}
|
|
|
|
// CreateIndex is an in-memory implementation of the Schemator CreateIndex method.
|
|
func (s *inMemSchemator) CreateIndex(ctx context.Context, name string, val []byte) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if idx, ok := s.schema[name]; ok {
|
|
// The current logic in pilosa doesn't allow us to return ErrIndexExists
|
|
// here, so for now we just update the Data value if the index already
|
|
// exists.
|
|
idx.Data = val
|
|
return nil
|
|
}
|
|
s.schema[name] = &Index{
|
|
Data: val,
|
|
Fields: make(map[string]*Field),
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// DeleteIndex is an in-memory implementation of the Schemator DeleteIndex method.
|
|
func (s *inMemSchemator) DeleteIndex(ctx context.Context, name string) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
delete(s.schema, name)
|
|
return nil
|
|
}
|
|
|
|
// Field is an in-memory implementation of the Schemator Field method.
|
|
func (s *inMemSchemator) Field(ctx context.Context, index, field string) ([]byte, error) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
idx, ok := s.schema[index]
|
|
if !ok {
|
|
return nil, ErrIndexDoesNotExist
|
|
}
|
|
fld, ok := idx.Fields[field]
|
|
if !ok {
|
|
return nil, ErrFieldDoesNotExist
|
|
}
|
|
return fld.Data, nil
|
|
}
|
|
|
|
// CreateField is an in-memory implementation of the Schemator CreateField method.
|
|
func (s *inMemSchemator) CreateField(ctx context.Context, index, field string, fieldVal []byte) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
idx, ok := s.schema[index]
|
|
if !ok {
|
|
return ErrIndexDoesNotExist
|
|
}
|
|
if fld, ok := idx.Fields[field]; ok {
|
|
// The current logic in pilosa doesn't allow us to return ErrFieldExists
|
|
// here, so for now we just update the Data value if the field already
|
|
// exists.
|
|
fld.Data = fieldVal
|
|
return nil
|
|
}
|
|
idx.Fields[field] = &Field{
|
|
Data: fieldVal,
|
|
Views: make(map[string]struct{}),
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *inMemSchemator) UpdateField(ctx context.Context, index, field string, fieldVal []byte) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
idx, ok := s.schema[index]
|
|
if !ok {
|
|
return ErrIndexDoesNotExist
|
|
}
|
|
if fld, ok := idx.Fields[field]; ok {
|
|
// The current logic in pilosa doesn't allow us to return ErrFieldExists
|
|
// here, so for now we just update the Data value if the field already
|
|
// exists.
|
|
fld.Data = fieldVal
|
|
return nil
|
|
} else {
|
|
return ErrFieldDoesNotExist
|
|
}
|
|
}
|
|
|
|
// DeleteField is an in-memory implementation of the Schemator DeleteField method.
|
|
func (s *inMemSchemator) DeleteField(ctx context.Context, index, field string) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
idx, ok := s.schema[index]
|
|
if !ok {
|
|
return ErrIndexDoesNotExist
|
|
}
|
|
delete(idx.Fields, field)
|
|
return nil
|
|
}
|
|
|
|
// View is an in-memory implementation of the Schemator View method.
|
|
func (s *inMemSchemator) View(ctx context.Context, index, field, view string) (bool, error) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
idx, ok := s.schema[index]
|
|
if !ok {
|
|
return false, ErrIndexDoesNotExist
|
|
}
|
|
fld, ok := idx.Fields[field]
|
|
if !ok {
|
|
return false, ErrFieldDoesNotExist
|
|
}
|
|
_, ok = fld.Views[view]
|
|
return ok, nil
|
|
}
|
|
|
|
// CreateView is an in-memory implementation of the Schemator CreateView method.
|
|
func (s *inMemSchemator) CreateView(ctx context.Context, index, field, view string) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
idx, ok := s.schema[index]
|
|
if !ok {
|
|
return ErrIndexDoesNotExist
|
|
}
|
|
fld, ok := idx.Fields[field]
|
|
if !ok {
|
|
return ErrFieldDoesNotExist
|
|
}
|
|
// The current logic in pilosa doesn't allow us to return ErrViewExists
|
|
// here, so for now we just update the value if the view already exists.
|
|
fld.Views[view] = struct{}{}
|
|
return nil
|
|
}
|
|
|
|
// DeleteView is an in-memory implementation of the Schemator DeleteView method.
|
|
func (s *inMemSchemator) DeleteView(ctx context.Context, index, field, view string) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
idx, ok := s.schema[index]
|
|
if !ok {
|
|
return ErrIndexDoesNotExist
|
|
}
|
|
fld, ok := idx.Fields[field]
|
|
if !ok {
|
|
return ErrFieldDoesNotExist
|
|
}
|
|
delete(fld.Views, view)
|
|
return nil
|
|
}
|
|
|
|
var InMemSharder Sharder = &inMemSharder{
|
|
shards: make(map[string][]byte),
|
|
}
|
|
|
|
type inMemSharder struct {
|
|
mu sync.RWMutex
|
|
shards map[string][]byte
|
|
}
|
|
|
|
func (s *inMemSharder) Shards(ctx context.Context, index, field string) ([][]byte, error) {
|
|
key := path.Join("/shard/", index, field)
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
|
|
b := s.shards[key]
|
|
if b == nil {
|
|
return nil, nil
|
|
}
|
|
return [][]byte{b}, nil
|
|
}
|
|
|
|
func (s *inMemSharder) SetShards(ctx context.Context, index, field string, shards []byte) error {
|
|
key := path.Join("/shard/", index, field)
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
|
|
s.shards[key] = make([]byte, len(shards))
|
|
copy(s.shards[key], shards)
|
|
return nil
|
|
}
|