mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-09 04:17:51 +00:00
Fixed some import issues resulting from syncing the private repo
This commit is contained in:
parent
e14cde7341
commit
45e405b980
38 changed files with 238 additions and 84 deletions
1
api.go
1
api.go
|
|
@ -23,7 +23,6 @@ import (
|
|||
"time"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/disco"
|
||||
"github.com/featurebasedb/featurebase/v3/ingest"
|
||||
"github.com/featurebasedb/featurebase/v3/rbf"
|
||||
|
||||
//"github.com/featurebasedb/featurebase/v3/pg"
|
||||
|
|
|
|||
|
|
@ -1,6 +1,5 @@
|
|||
// Copyright 2022 Molecula Corp. (DBA FeatureBase).
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
package client
|
||||
// Package batch provides tooling to prepare batches of records for ingest.
|
||||
package batch
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
|
|
@ -11,8 +10,9 @@ import (
|
|||
"time"
|
||||
|
||||
featurebase "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/client/egpool"
|
||||
"github.com/featurebasedb/featurebase/v3/batch/egpool"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
"github.com/featurebasedb/featurebase/v3/pql"
|
||||
"github.com/featurebasedb/featurebase/v3/roaring"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
|
@ -716,7 +716,7 @@ func (b *Batch) Add(rec Row) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// ErrBatchNowFull — similar to io.EOF — is a marker error to notify the user of
|
||||
// ErrBatchNowFull — similar to io.EOF — is a marker error to notify the user of
|
||||
// a batch that it is time to call Import.
|
||||
var ErrBatchNowFull = errors.New("batch is now full - you cannot add any more records (though the one you just added was accepted)")
|
||||
|
||||
|
|
|
|||
|
|
@ -10,9 +10,9 @@ import (
|
|||
"testing"
|
||||
"time"
|
||||
|
||||
featurebase "github.com/molecula/featurebase/v3"
|
||||
"github.com/molecula/featurebase/v3/client"
|
||||
"github.com/molecula/featurebase/v3/pql"
|
||||
featurebase "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/client"
|
||||
"github.com/featurebasedb/featurebase/v3/pql"
|
||||
"github.com/stretchr/testify/assert"
|
||||
|
||||
"github.com/pkg/errors"
|
||||
|
|
|
|||
|
|
@ -3,8 +3,8 @@ package batch
|
|||
import (
|
||||
"time"
|
||||
|
||||
featurebase "github.com/molecula/featurebase/v3"
|
||||
"github.com/molecula/featurebase/v3/errors"
|
||||
featurebase "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/errors"
|
||||
)
|
||||
|
||||
var (
|
||||
|
|
|
|||
|
|
@ -6,7 +6,7 @@ import (
|
|||
"errors"
|
||||
"testing"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/client/egpool"
|
||||
"github.com/featurebasedb/featurebase/v3/batch/egpool"
|
||||
)
|
||||
|
||||
func TestEGPool(t *testing.T) {
|
||||
|
|
|
|||
|
|
@ -5,10 +5,10 @@ import (
|
|||
"time"
|
||||
|
||||
"github.com/golang/protobuf/proto" //nolint:staticcheck
|
||||
featurebase "github.com/molecula/featurebase/v3"
|
||||
featurebaseproto "github.com/molecula/featurebase/v3/encoding/proto"
|
||||
"github.com/molecula/featurebase/v3/pb"
|
||||
"github.com/molecula/featurebase/v3/roaring"
|
||||
featurebase "github.com/featurebasedb/featurebase/v3"
|
||||
featurebaseproto "github.com/featurebasedb/featurebase/v3/encoding/proto"
|
||||
"github.com/featurebasedb/featurebase/v3/pb"
|
||||
"github.com/featurebasedb/featurebase/v3/roaring"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
// Copyright 2022 Molecula Corp. (DBA FeatureBase).
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
package client
|
||||
package batch
|
||||
|
||||
const (
|
||||
// MetricBatchImportDurationSeconds records the full time of the
|
||||
|
|
|
|||
|
|
@ -3,8 +3,8 @@ package client
|
|||
import (
|
||||
"context"
|
||||
|
||||
featurebase "github.com/molecula/featurebase/v3"
|
||||
"github.com/molecula/featurebase/v3/errors"
|
||||
featurebase "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/errors"
|
||||
)
|
||||
|
||||
var _ featurebase.SchemaAPI = &schemaAPI{}
|
||||
|
|
|
|||
|
|
@ -4,8 +4,8 @@ import (
|
|||
"context"
|
||||
"time"
|
||||
|
||||
featurebase "github.com/molecula/featurebase/v3"
|
||||
"github.com/molecula/featurebase/v3/roaring"
|
||||
featurebase "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/roaring"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -9,7 +9,6 @@ import (
|
|||
"time"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/disco"
|
||||
"github.com/featurebasedb/featurebase/v3/ingest"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
"github.com/featurebasedb/featurebase/v3/roaring"
|
||||
"github.com/pkg/errors"
|
||||
|
|
|
|||
|
|
@ -1,5 +1,4 @@
|
|||
// Copyright 2022 Molecula Corp. (DBA FeatureBase).
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// Copyright 2021 Molecula Corp. All rights reserved.
|
||||
package main
|
||||
|
||||
import (
|
||||
|
|
@ -8,16 +7,17 @@ import (
|
|||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
gohttp "net/http"
|
||||
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/encoding/proto"
|
||||
pnet "github.com/featurebasedb/featurebase/v3/net"
|
||||
"github.com/featurebasedb/featurebase/v3/vprint"
|
||||
|
||||
"os"
|
||||
"strconv"
|
||||
"strings"
|
||||
)
|
||||
|
||||
func UploadTar(srcFile string, client *pilosa.InternalClient) error {
|
||||
|
|
|
|||
|
|
@ -10,7 +10,6 @@ import (
|
|||
"github.com/gogo/protobuf/proto"
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/disco"
|
||||
"github.com/featurebasedb/featurebase/v3/ingest"
|
||||
pnet "github.com/featurebasedb/featurebase/v3/net"
|
||||
"github.com/featurebasedb/featurebase/v3/pb"
|
||||
"github.com/featurebasedb/featurebase/v3/pql"
|
||||
|
|
|
|||
|
|
@ -7,7 +7,6 @@ import (
|
|||
"testing"
|
||||
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/ingest"
|
||||
"github.com/featurebasedb/featurebase/v3/pb"
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -19,10 +19,7 @@ import (
|
|||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/cespare/xxhash"
|
||||
"github.com/featurebasedb/featurebase/v3/disco"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
pnet "github.com/featurebasedb/featurebase/v3/net"
|
||||
"github.com/featurebasedb/featurebase/v3/pb"
|
||||
"github.com/featurebasedb/featurebase/v3/pql"
|
||||
"github.com/featurebasedb/featurebase/v3/roaring"
|
||||
|
|
|
|||
|
|
@ -1,5 +1,4 @@
|
|||
// Copyright 2022 Molecula Corp. (DBA FeatureBase).
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// Copyright 2021 Molecula Corp. All rights reserved.
|
||||
package pilosa
|
||||
|
||||
import (
|
||||
|
|
@ -7,7 +6,7 @@ import (
|
|||
"math/bits"
|
||||
"time"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/ingest"
|
||||
"github.com/featurebasedb/featurebase/v3/shardwidth"
|
||||
"github.com/featurebasedb/featurebase/v3/tracing"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@ import (
|
|||
"math/rand"
|
||||
"testing"
|
||||
|
||||
pilosa "github.com/molecula/featurebase/v3"
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
)
|
||||
|
||||
func TestSortToShards(t *testing.T) {
|
||||
|
|
|
|||
|
|
@ -31,7 +31,6 @@ import (
|
|||
"github.com/featurebasedb/featurebase/v3/authn"
|
||||
"github.com/featurebasedb/featurebase/v3/authz"
|
||||
"github.com/featurebasedb/featurebase/v3/disco"
|
||||
"github.com/featurebasedb/featurebase/v3/ingest"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
"github.com/featurebasedb/featurebase/v3/monitor"
|
||||
"github.com/featurebasedb/featurebase/v3/pql"
|
||||
|
|
|
|||
|
|
@ -21,7 +21,6 @@ import (
|
|||
|
||||
"github.com/featurebasedb/featurebase/v3/authn"
|
||||
"github.com/featurebasedb/featurebase/v3/disco"
|
||||
"github.com/featurebasedb/featurebase/v3/ingest"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
pnet "github.com/featurebasedb/featurebase/v3/net"
|
||||
"github.com/featurebasedb/featurebase/v3/tracing"
|
||||
|
|
|
|||
|
|
@ -5,7 +5,6 @@ package pilosa
|
|||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/roaring"
|
||||
)
|
||||
|
||||
// iterator is an interface for looping over row/column pairs.
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@ import (
|
|||
"sort"
|
||||
"strings"
|
||||
|
||||
"github.com/molecula/featurebase/v3/roaring"
|
||||
"github.com/featurebasedb/featurebase/v3/roaring"
|
||||
)
|
||||
|
||||
// QueryContext represents the lifespan of a query or similar thing which
|
||||
|
|
|
|||
|
|
@ -11,7 +11,7 @@ import (
|
|||
"testing"
|
||||
"time"
|
||||
|
||||
rbfcfg "github.com/molecula/featurebase/v3/rbf/cfg"
|
||||
rbfcfg "github.com/featurebasedb/featurebase/v3/rbf/cfg"
|
||||
"golang.org/x/sync/errgroup"
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -10,9 +10,9 @@ import (
|
|||
"runtime"
|
||||
"sync"
|
||||
|
||||
"github.com/molecula/featurebase/v3/rbf"
|
||||
rbfcfg "github.com/molecula/featurebase/v3/rbf/cfg"
|
||||
"github.com/molecula/featurebase/v3/roaring"
|
||||
"github.com/featurebasedb/featurebase/v3/rbf"
|
||||
rbfcfg "github.com/featurebasedb/featurebase/v3/rbf/cfg"
|
||||
"github.com/featurebasedb/featurebase/v3/roaring"
|
||||
)
|
||||
|
||||
// rbfDBQueryContexts represents an actual backend DB, and a map of the
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@ import (
|
|||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/molecula/featurebase/v3/roaring"
|
||||
"github.com/featurebasedb/featurebase/v3/roaring"
|
||||
)
|
||||
|
||||
func TestRbfWrite(t *testing.T) {
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@ import (
|
|||
"sync/atomic"
|
||||
"testing"
|
||||
|
||||
"github.com/molecula/featurebase/v3/roaring"
|
||||
"github.com/featurebasedb/featurebase/v3/roaring"
|
||||
)
|
||||
|
||||
type oopsieWrapper struct {
|
||||
|
|
|
|||
1
rbf.go
1
rbf.go
|
|
@ -16,7 +16,6 @@ import (
|
|||
txkey "github.com/featurebasedb/featurebase/v3/short_txkey"
|
||||
"github.com/featurebasedb/featurebase/v3/storage"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/vprint"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
|
|
|
|||
178
server/server.go
178
server/server.go
|
|
@ -1,5 +1,4 @@
|
|||
// Copyright 2022 Molecula Corp. (DBA FeatureBase).
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// Copyright 2021 Molecula Corp. All rights reserved.
|
||||
//
|
||||
// Package server contains the `pilosa server` subcommand which runs Pilosa
|
||||
// itself. The purpose of this package is to define an easily tested Command
|
||||
|
|
@ -16,6 +15,7 @@ import (
|
|||
"log"
|
||||
"math/rand"
|
||||
"net"
|
||||
"net/http"
|
||||
"os"
|
||||
"os/signal"
|
||||
"path/filepath"
|
||||
|
|
@ -26,12 +26,16 @@ import (
|
|||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/dax"
|
||||
"golang.org/x/sync/errgroup"
|
||||
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/authn"
|
||||
"github.com/featurebasedb/featurebase/v3/authz"
|
||||
"github.com/featurebasedb/featurebase/v3/batch"
|
||||
"github.com/featurebasedb/featurebase/v3/boltdb"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/computer"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/computer/alpha"
|
||||
"github.com/featurebasedb/featurebase/v3/encoding/proto"
|
||||
petcd "github.com/featurebasedb/featurebase/v3/etcd"
|
||||
"github.com/featurebasedb/featurebase/v3/gcnotify"
|
||||
|
|
@ -75,7 +79,12 @@ type Command struct {
|
|||
logger loggerLogger
|
||||
queryLogger loggerLogger
|
||||
|
||||
mds pilosa.MDS
|
||||
writeLogger pilosa.WriteLogger
|
||||
snapshotter pilosa.Snapshotter
|
||||
|
||||
Handler pilosa.HandlerI
|
||||
httpHandler http.Handler
|
||||
grpcServer *grpcServer
|
||||
grpcLn net.Listener
|
||||
API *pilosa.API
|
||||
|
|
@ -86,6 +95,10 @@ type Command struct {
|
|||
|
||||
serverOptions []pilosa.ServerOption
|
||||
auth *authn.Auth
|
||||
|
||||
// isComputeNode is set to true if this node is running as a DAX compute
|
||||
// node.
|
||||
isComputeNode bool
|
||||
}
|
||||
|
||||
type CommandOption func(c *Command) error
|
||||
|
|
@ -111,6 +124,8 @@ func OptCommandConfig(config *Config) CommandOption {
|
|||
c.Config.Etcd = config.Etcd
|
||||
c.Config.Auth = config.Auth
|
||||
c.Config.TLS = config.TLS
|
||||
c.Config.MDSAddress = config.MDSAddress
|
||||
c.Config.WriteLogger = config.WriteLogger
|
||||
return nil
|
||||
}
|
||||
c.Config = config
|
||||
|
|
@ -118,6 +133,41 @@ func OptCommandConfig(config *Config) CommandOption {
|
|||
}
|
||||
}
|
||||
|
||||
// OptCommandSetConfig was added because OptCommandConfig only sets a small
|
||||
// sub-set of the config options (it doesn't seem to be used for anything but
|
||||
// tests). We need a functional option which sets the full Config.
|
||||
func OptCommandSetConfig(config *Config) CommandOption {
|
||||
return func(c *Command) error {
|
||||
defer c.Config.MustValidate()
|
||||
c.Config = config
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// OptCommandInjections injects the interface implementations.
|
||||
func OptCommandInjections(inj Injections) CommandOption {
|
||||
return func(c *Command) error {
|
||||
if inj.MDS != nil {
|
||||
c.mds = inj.MDS
|
||||
}
|
||||
if inj.WriteLogger != nil {
|
||||
c.writeLogger = inj.WriteLogger
|
||||
}
|
||||
if inj.Snapshotter != nil {
|
||||
c.snapshotter = inj.Snapshotter
|
||||
}
|
||||
c.isComputeNode = inj.IsComputeNode
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
type Injections struct {
|
||||
MDS pilosa.MDS
|
||||
WriteLogger pilosa.WriteLogger
|
||||
Snapshotter pilosa.Snapshotter
|
||||
IsComputeNode bool
|
||||
}
|
||||
|
||||
// NewCommand returns a new instance of Main.
|
||||
func NewCommand(stdin io.Reader, stdout, stderr io.Writer, opts ...CommandOption) *Command {
|
||||
c := &Command{
|
||||
|
|
@ -183,7 +233,7 @@ func (m *Command) doSetupResourceLimits() error {
|
|||
return fmt.Errorf("checking open file limit: %w", err)
|
||||
} else {
|
||||
if oldLimit.Cur != targetFileLimit {
|
||||
m.logger.Warnf("tried to set open file limit to %d, but it is %d; see https://docs.featurebase.com/reference/hostsystem#operating-system-configuration", targetFileLimit, oldLimit.Cur)
|
||||
m.logger.Warnf("tried to set open file limit to %d, but it is %d; see https://docs.featurebasedb.cloud/reference/hostsystem#operating-system-configuration", targetFileLimit, oldLimit.Cur)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -215,12 +265,70 @@ func (m *Command) setupResourceLimits() error {
|
|||
return setupResourceLimitsErr
|
||||
}
|
||||
|
||||
// StartNoServe starts the pilosa server, but doesn't serve on the http handler.
|
||||
func (m *Command) StartNoServe() (err error) {
|
||||
// Seed random number generator
|
||||
rand.Seed(time.Now().UTC().UnixNano())
|
||||
|
||||
// setupServer
|
||||
err = m.setupServer()
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "setting up server")
|
||||
}
|
||||
err = m.setupResourceLimits()
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "setting resource limits")
|
||||
}
|
||||
|
||||
// Initialize server.
|
||||
if err = m.Server.Open(); err != nil {
|
||||
return errors.Wrap(err, "opening server")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Register registers the node with the MDS service using whatever MDS
|
||||
// implementation was injected during setup.
|
||||
func (m *Command) Register() (err error) {
|
||||
if m.mds == nil {
|
||||
return errors.New("no MDS implementation with which to register")
|
||||
}
|
||||
|
||||
node := &dax.Node{
|
||||
Address: dax.Address(m.Config.Advertise),
|
||||
RoleTypes: []dax.RoleType{
|
||||
dax.RoleTypeCompute,
|
||||
dax.RoleTypeTranslate,
|
||||
},
|
||||
}
|
||||
return m.mds.RegisterNode(context.Background(), node)
|
||||
}
|
||||
|
||||
// CheckIn is called periodically to check in with the MDS service using
|
||||
// whatever MDS implementation was injected during setup.
|
||||
func (m *Command) CheckIn() (err error) {
|
||||
if m.mds == nil {
|
||||
return errors.New("no MDS implementation with which to check-in")
|
||||
}
|
||||
|
||||
node := &dax.Node{
|
||||
Address: dax.Address(m.Config.Advertise),
|
||||
RoleTypes: []dax.RoleType{
|
||||
dax.RoleTypeCompute,
|
||||
dax.RoleTypeTranslate,
|
||||
},
|
||||
}
|
||||
return m.mds.CheckInNode(context.Background(), node)
|
||||
}
|
||||
|
||||
// Start starts the pilosa server - it returns once the server is running.
|
||||
func (m *Command) Start() (err error) {
|
||||
// Seed random number generator
|
||||
rand.Seed(time.Now().UTC().UnixNano())
|
||||
// SetupServer
|
||||
err = m.SetupServer()
|
||||
|
||||
// setupServer
|
||||
err = m.setupServer()
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "setting up server")
|
||||
}
|
||||
|
|
@ -258,8 +366,8 @@ func (m *Command) UpAndDown() (err error) {
|
|||
// Seed random number generator
|
||||
rand.Seed(time.Now().UTC().UnixNano())
|
||||
|
||||
// SetupServer
|
||||
err = m.SetupServer()
|
||||
// setupServer
|
||||
err = m.setupServer()
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "setting up server")
|
||||
}
|
||||
|
|
@ -300,8 +408,8 @@ func (m *Command) Wait() error {
|
|||
}
|
||||
}
|
||||
|
||||
// SetupServer uses the cluster configuration to set up this server.
|
||||
func (m *Command) SetupServer() error {
|
||||
// setupServer uses the cluster configuration to set up this server.
|
||||
func (m *Command) setupServer() error {
|
||||
runtime.SetBlockProfileRate(m.Config.Profile.BlockRate)
|
||||
runtime.SetMutexProfileFraction(m.Config.Profile.MutexFraction)
|
||||
|
||||
|
|
@ -368,9 +476,13 @@ func (m *Command) SetupServer() error {
|
|||
return errors.Wrap(err, "new stats client")
|
||||
}
|
||||
|
||||
m.ln, err = getListener(*uri, m.tlsConfig)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "getting listener")
|
||||
if m.Config.Listener == nil {
|
||||
m.ln, err = getListener(*uri, m.tlsConfig)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "getting listener")
|
||||
}
|
||||
} else {
|
||||
m.ln = m.Config.Listener
|
||||
}
|
||||
|
||||
// If port is 0, get auto-allocated port from listener
|
||||
|
|
@ -439,6 +551,21 @@ func (m *Command) SetupServer() error {
|
|||
m.Config.Etcd.Dir = filepath.Join(path, pilosa.DiscoDir)
|
||||
}
|
||||
|
||||
// WriteLogger setup.
|
||||
var wlw computer.WriteLogWriter = computer.NewNopWriteLogWriter()
|
||||
var wlr computer.WriteLogReader = computer.NewNopWriteLogReader()
|
||||
if m.writeLogger != nil {
|
||||
alphaWriteLog := alpha.NewAlphaWriteLog(m.writeLogger)
|
||||
wlr = alphaWriteLog
|
||||
wlw = alphaWriteLog
|
||||
}
|
||||
|
||||
// Snapshotter setup.
|
||||
var snap computer.SnapshotReadWriter = computer.NewNopSnapshotReadWriter()
|
||||
if m.snapshotter != nil {
|
||||
snap = alpha.NewAlphaSnapshot(m.snapshotter)
|
||||
}
|
||||
|
||||
m.Config.Etcd.Id = m.Config.Name // TODO(twg) rethink this
|
||||
e := petcd.NewEtcd(m.Config.Etcd, m.logger, m.Config.Cluster.ReplicaN, version)
|
||||
|
||||
|
|
@ -477,6 +604,9 @@ func (m *Command) SetupServer() error {
|
|||
pilosa.OptServerPartitionAssigner(m.Config.Cluster.PartitionToNodeAssignment),
|
||||
pilosa.OptServerDisCo(e, e, e, e),
|
||||
pilosa.OptServerExecutionPlannerFn(executionPlannerFn),
|
||||
pilosa.OptServerWriteLogReader(wlr),
|
||||
pilosa.OptServerWriteLogWriter(wlw),
|
||||
pilosa.OptServerSnapshotReadWriter(snap),
|
||||
}
|
||||
|
||||
if m.Config.LookupDBDSN != "" {
|
||||
|
|
@ -492,7 +622,6 @@ func (m *Command) SetupServer() error {
|
|||
}
|
||||
|
||||
m.Server, err = pilosa.NewServer(serverOptions...)
|
||||
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "new server")
|
||||
}
|
||||
|
|
@ -500,6 +629,11 @@ func (m *Command) SetupServer() error {
|
|||
m.API, err = pilosa.NewAPI(
|
||||
pilosa.OptAPIServer(m.Server),
|
||||
pilosa.OptAPIImportWorkerPoolSize(m.Config.ImportWorkerPoolSize),
|
||||
pilosa.OptAPIWriteLogReader(wlr),
|
||||
pilosa.OptAPIWriteLogWriter(wlw),
|
||||
pilosa.OptAPISnapshotter(snap),
|
||||
pilosa.OptAPIDirectiveWorkerPoolSize(m.Config.DirectiveWorkerPoolSize),
|
||||
pilosa.OptAPIIsComputeNode(m.isComputeNode),
|
||||
)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "new api")
|
||||
|
|
@ -559,7 +693,7 @@ func (m *Command) SetupServer() error {
|
|||
return errors.Wrap(err, "getting grpcServer")
|
||||
}
|
||||
|
||||
m.Handler, err = pilosa.NewHandler(
|
||||
hndlr, err := pilosa.NewHandler(
|
||||
pilosa.OptHandlerAllowedOrigins(m.Config.Handler.AllowedOrigins),
|
||||
pilosa.OptHandlerAPI(m.API),
|
||||
pilosa.OptHandlerLogger(m.logger),
|
||||
|
|
@ -574,7 +708,21 @@ func (m *Command) SetupServer() error {
|
|||
pilosa.OptHandlerRoaringSerializer(proto.RoaringSerializer),
|
||||
pilosa.OptHandlerSQLEnabled(m.Config.SQL.EndpointEnabled),
|
||||
)
|
||||
return errors.Wrap(err, "new handler")
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "new handler")
|
||||
}
|
||||
|
||||
m.httpHandler = hndlr
|
||||
m.Handler = hndlr
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// HTTPHandler was added for the case where we want to get the full
|
||||
// http.Handler, and not just those methods which satisfy the pilosa.HandlerI
|
||||
// interface.
|
||||
func (m *Command) HTTPHandler() http.Handler {
|
||||
return m.httpHandler
|
||||
}
|
||||
|
||||
// setupLogger sets up the logger based on the configuration.
|
||||
|
|
|
|||
|
|
@ -1,5 +1,4 @@
|
|||
// Copyright 2022 Molecula Corp. (DBA FeatureBase).
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// Copyright 2021 Molecula Corp. All rights reserved.
|
||||
package sql
|
||||
|
||||
import (
|
||||
|
|
@ -7,6 +6,7 @@ import (
|
|||
"fmt"
|
||||
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/authn"
|
||||
pproto "github.com/featurebasedb/featurebase/v3/proto"
|
||||
"github.com/pkg/errors"
|
||||
"google.golang.org/grpc/codes"
|
||||
|
|
|
|||
|
|
@ -7,6 +7,8 @@ import (
|
|||
"encoding/json"
|
||||
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/batch"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3/parser"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3/planner/types"
|
||||
|
|
|
|||
|
|
@ -5,9 +5,9 @@ package planner
|
|||
import (
|
||||
"context"
|
||||
|
||||
pilosa "github.com/molecula/featurebase/v3"
|
||||
"github.com/molecula/featurebase/v3/sql3"
|
||||
"github.com/molecula/featurebase/v3/sql3/parser"
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3/parser"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -3,6 +3,8 @@
|
|||
package planner
|
||||
|
||||
import (
|
||||
"strings"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/sql3"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3/parser"
|
||||
)
|
||||
|
|
|
|||
|
|
@ -15,7 +15,9 @@ import (
|
|||
"strings"
|
||||
"time"
|
||||
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/PaesslerAG/gval"
|
||||
"github.com/PaesslerAG/jsonpath"
|
||||
"github.com/featurebasedb/featurebase/v3/pql"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3/parser"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3/planner/types"
|
||||
|
|
|
|||
|
|
@ -9,10 +9,10 @@ import (
|
|||
"time"
|
||||
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/pql"
|
||||
fbbatch "github.com/featurebasedb/featurebase/v3/batch"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3/parser"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3/planner/types"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
// PlanOpInsert plan operator to handle INSERT.
|
||||
|
|
|
|||
|
|
@ -6,10 +6,10 @@ import (
|
|||
"context"
|
||||
"fmt"
|
||||
|
||||
pilosa "github.com/molecula/featurebase/v3"
|
||||
"github.com/molecula/featurebase/v3/sql3"
|
||||
"github.com/molecula/featurebase/v3/sql3/parser"
|
||||
"github.com/molecula/featurebase/v3/sql3/planner/types"
|
||||
pilosa "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3/parser"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3/planner/types"
|
||||
)
|
||||
|
||||
//fb_exec_requests
|
||||
|
|
|
|||
|
|
@ -5,9 +5,9 @@ import (
|
|||
"github.com/featurebasedb/featurebase/v3/sql3/parser"
|
||||
)
|
||||
|
||||
//aggregate function tests
|
||||
var countTests = tableTest{
|
||||
table: tbl(
|
||||
// aggregate function tests
|
||||
var countTests = TableTest{
|
||||
Table: tbl(
|
||||
"count_test",
|
||||
srcHdrs(
|
||||
srcHdr("_id", fldTypeID),
|
||||
|
|
@ -51,7 +51,20 @@ var countTests = tableTest{
|
|||
Compare: CompareExactUnordered,
|
||||
},
|
||||
{
|
||||
sqls: sqls(
|
||||
SQLs: sqls(
|
||||
"SELECT COUNT(i1) as a, COUNT(i2) as b FROM count_test",
|
||||
),
|
||||
ExpHdrs: hdrs(
|
||||
hdr("a", fldTypeInt),
|
||||
hdr("b", fldTypeInt),
|
||||
),
|
||||
ExpRows: rows(
|
||||
row(int64(6), int64(2)),
|
||||
),
|
||||
Compare: CompareExactUnordered,
|
||||
},
|
||||
{
|
||||
SQLs: sqls(
|
||||
"SELECT COUNT(*) + 10 - 11 * 2 AS count_rows FROM count_test",
|
||||
),
|
||||
ExpHdrs: hdrs(
|
||||
|
|
|
|||
|
|
@ -3,7 +3,7 @@ package defs
|
|||
import (
|
||||
"time"
|
||||
|
||||
"github.com/molecula/featurebase/v3/pql"
|
||||
"github.com/featurebasedb/featurebase/v3/pql"
|
||||
)
|
||||
|
||||
// INT bin op tests
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package defs
|
||||
|
||||
import "github.com/molecula/featurebase/v3/pql"
|
||||
import "github.com/featurebasedb/featurebase/v3/pql"
|
||||
|
||||
var unkeyed = TableTest{
|
||||
name: "unkeyed",
|
||||
|
|
|
|||
|
|
@ -8,8 +8,8 @@ import (
|
|||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/molecula/featurebase/v3/sql3/parser"
|
||||
planner_types "github.com/molecula/featurebase/v3/sql3/planner/types"
|
||||
"github.com/featurebasedb/featurebase/v3/sql3/parser"
|
||||
planner_types "github.com/featurebasedb/featurebase/v3/sql3/planner/types"
|
||||
)
|
||||
|
||||
type fldType parser.ExprDataType
|
||||
|
|
|
|||
|
|
@ -11,7 +11,6 @@ import (
|
|||
"sync"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/disco"
|
||||
"github.com/featurebasedb/featurebase/v3/ingest"
|
||||
"github.com/featurebasedb/featurebase/v3/roaring"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue