diff --git a/bolt.go b/bolt.go index b8a3f25be..6e59b9af5 100644 --- a/bolt.go +++ b/bolt.go @@ -29,9 +29,8 @@ import ( "time" "github.com/pilosa/pilosa/v2/hash" - "github.com/pilosa/pilosa/v2/rbf" - rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg" "github.com/pilosa/pilosa/v2/roaring" + "github.com/pilosa/pilosa/v2/storage" // On Bolt only, we still use the long txkey, because // this allows Max() to work readily. @@ -130,7 +129,7 @@ func boltPath(path string) string { // if one does not exist for its bpath. Otherwise it returns // the existing instance. This insures only one boltDB // per bpath in this pilosa node. -func (r *boltRegistrar) OpenDBWrapper(path0 string, doAllocZero bool, rbfcfg *rbfcfg.Config) (DBWrapper, error) { +func (r *boltRegistrar) OpenDBWrapper(path0 string, doAllocZero bool, cfg *storage.Config) (DBWrapper, error) { path := boltPath(path0) r.mu.Lock() @@ -171,7 +170,7 @@ func (r *boltRegistrar) OpenDBWrapper(path0 string, doAllocZero bool, rbfcfg *rb // re-sync during recovery. // NoFreelistSync bool - if rbfcfg != nil && !rbfcfg.FsyncEnabled { + if cfg != nil && !cfg.FsyncEnabled { db.NoSync = true db.NoFreelistSync = true } else { @@ -479,7 +478,7 @@ func (tx *BoltTx) Type() string { } func (tx *BoltTx) UseRowCache() bool { - return rbf.EnableRowCache() + return storage.EnableRowCache() } // Pointer gives us a memory address for the underlying transaction for debugging. diff --git a/ctl/server.go b/ctl/server.go index 65f070c26..07fdf887a 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -20,6 +20,7 @@ import ( "github.com/pilosa/pilosa/v2" "github.com/pilosa/pilosa/v2/server" + "github.com/pilosa/pilosa/v2/storage" "github.com/spf13/cobra" ) @@ -107,6 +108,15 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) { // cannot detect and honor the PILOSA_TXSRC env var over-ride. flags.StringVarP(&srv.Config.Txsrc, "txsrc", "", "", fmt.Sprintf("transaction/storage to use: one of roaring, rbf, bolt, or a blue-green setup: rbf_roaring, roaring_rbf, bolt_roaring, roaring_bolt, bolt_rbf, etc. The default is: %v. The env var PILOSA_TXSRC is over-ridden by --txsrc option on the command line.", pilosa.DefaultTxsrc)) + // Storage + // Note: the default for --storage.backend must be kept "" empty string. + // Otherwise we cannot detect and honor the PILOSA_STORAGE_BACKEND env var + // over-ride. + // TODO: the comment above was carried over from the PILOSA_TXSRC flag, but + // we should confirm that this still applies. + flags.StringVar(&srv.Config.Storage.Backend, "storage.backend", storage.DefaultBackend, fmt.Sprintf("transaction/storage to use: one of roaring, rbf, bolt, or a blue-green setup: rbf_roaring, roaring_rbf, bolt_roaring, roaring_bolt, bolt_rbf, etc. The default is: %v. The env var PILOSA_STORAGE_BACKEND is over-ridden by --storage.backend option on the command line.", storage.DefaultBackend)) + flags.BoolVar(&srv.Config.Storage.FsyncEnabled, "storage.fsync", true, "enable fsync fully safe flush-to-disk") + // RowcacheOn flags.BoolVarP((&srv.Config.RowcacheOn), "rowcache-on", "", srv.Config.RowcacheOn, "turn on the rowcache for all backends (may speed some queries)") diff --git a/dbshard.go b/dbshard.go index 55d31ed8a..d08f06b87 100644 --- a/dbshard.go +++ b/dbshard.go @@ -25,6 +25,8 @@ import ( rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg" txkey "github.com/pilosa/pilosa/v2/short_txkey" + "github.com/pilosa/pilosa/v2/storage" + //txkey "github.com/pilosa/pilosa/v2/txkey" "github.com/pkg/errors" ) @@ -61,7 +63,7 @@ type DBWrapper interface { } type DBRegistry interface { - OpenDBWrapper(path string, doAllocZero bool, rbfcfg *rbfcfg.Config) (DBWrapper, error) + OpenDBWrapper(path string, doAllocZero bool, cfg *storage.Config) (DBWrapper, error) } type DBShard struct { @@ -243,7 +245,8 @@ type DBPerShard struct { isBlueGreen bool - RBFConfig *rbfcfg.Config + StorageConfig *storage.Config + RBFConfig *rbfcfg.Config } func newIndex2Shards() (r map[txtype]map[string]*shardSet) { @@ -399,9 +402,8 @@ func (per *DBPerShard) LoadExistingDBs() (err error) { } func (txf *TxFactory) NewDBPerShard(types []txtype, holderDir string, holder *Holder) (d *DBPerShard) { - - if holder.cfg == nil || holder.cfg.RBFConfig == nil { - panic("must have holder.cfg.RBFConfig set here") + if holder.cfg == nil || holder.cfg.RBFConfig == nil || holder.cfg.StorageConfig == nil { + panic("must have holder.cfg.RBFConfig and holder.cfg.StorageConfig set here") } useOpenList := 0 @@ -421,17 +423,18 @@ func (txf *TxFactory) NewDBPerShard(types []txtype, holderDir string, holder *Ho } d = &DBPerShard{ - types: types, - HolderDir: holderDir, - holder: holder, - dbh: NewDBHolder(), - Flatmap: make(map[flatkey]*DBShard), - txf: txf, - useOpenList: useOpenList, - hasRoaring: hasRoaring, - isBlueGreen: len(types) > 1, - index2shards: newIndex2Shards(), - RBFConfig: holder.cfg.RBFConfig, + types: types, + HolderDir: holderDir, + holder: holder, + dbh: NewDBHolder(), + Flatmap: make(map[flatkey]*DBShard), + txf: txf, + useOpenList: useOpenList, + hasRoaring: hasRoaring, + isBlueGreen: len(types) > 1, + index2shards: newIndex2Shards(), + StorageConfig: holder.cfg.StorageConfig, + RBFConfig: holder.cfg.RBFConfig, } return } @@ -645,13 +648,14 @@ func (per *DBPerShard) unprotectedGetDBShard(index string, shard uint64, idx *In registry = globalRoaringReg case rbfTxn: registry = globalRbfDBReg + registry.(*rbfDBRegistrar).SetRBFConfig(per.RBFConfig) case boltTxn: registry = globalBoltReg default: panic(fmt.Sprintf("unknown txtyp: '%v'", ty)) } path := dbs.pathForType(ty) - w, err := registry.OpenDBWrapper(path, DetectMemAccessPastTx, per.RBFConfig) + w, err := registry.OpenDBWrapper(path, DetectMemAccessPastTx, per.StorageConfig) panicOn(err) h := idx.Holder() w.SetHolder(h) diff --git a/debugstats/stats.go b/debugstats/stats.go index 91005482a..edb9a3600 100644 --- a/debugstats/stats.go +++ b/debugstats/stats.go @@ -17,7 +17,6 @@ package debugstats import ( "fmt" "math" - //"os" "runtime" "sort" "sync" @@ -67,7 +66,6 @@ func (p SortByTot) Swap(i, j int) { } func (c *CallStats) Report(title string) (r string) { - //txsrc := os.Getenv("PILOSA_TXSRC") r = fmt.Sprintf("CallStats: (%v)\n", title) c.mu.Lock() defer c.mu.Unlock() diff --git a/executor_test.go b/executor_test.go index 6db994817..906d0d276 100644 --- a/executor_test.go +++ b/executor_test.go @@ -41,6 +41,7 @@ import ( "github.com/pilosa/pilosa/v2/pql" "github.com/pilosa/pilosa/v2/proto" "github.com/pilosa/pilosa/v2/server" + "github.com/pilosa/pilosa/v2/storage" "github.com/pilosa/pilosa/v2/test" "github.com/pilosa/pilosa/v2/testhook" "github.com/pkg/errors" @@ -6689,7 +6690,11 @@ func TestTimelessClearRegression(t *testing.T) { } func TestMissingKeyRegression(t *testing.T) { - c := test.MustRunCluster(t, 1, []server.CommandOption{server.OptCommandServerOptions(pilosa.OptServerTxsrc("roaring"))}) + c := test.MustRunCluster(t, 1, []server.CommandOption{server.OptCommandServerOptions( + pilosa.OptServerStorageConfig(&storage.Config{ + Backend: "roaring", + FsyncEnabled: true, + }))}) defer c.Close() c.CreateField(t, "i", pilosa.IndexOptions{Keys: true, TrackExistence: true}, "f", pilosa.OptFieldKeys()) diff --git a/holder.go b/holder.go index 2cd6daf73..21b4ada8e 100644 --- a/holder.go +++ b/holder.go @@ -31,10 +31,10 @@ import ( "time" "github.com/pilosa/pilosa/v2/logger" - "github.com/pilosa/pilosa/v2/rbf" rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg" "github.com/pilosa/pilosa/v2/roaring" "github.com/pilosa/pilosa/v2/stats" + "github.com/pilosa/pilosa/v2/storage" "github.com/pilosa/pilosa/v2/testhook" "github.com/pilosa/pilosa/v2/topology" "github.com/pilosa/pilosa/v2/tracing" @@ -150,7 +150,7 @@ type HolderOpts struct { Inspect bool // Txsrc controls the tx/storage engine we instatiate. Set by - // server.go OptServerTxsrc + // server.go OptServerStorageConfig Txsrc string // RowcacheOn, if true, turns on the row cache for all storage backends. @@ -214,9 +214,9 @@ type HolderConfig struct { StatsClient stats.StatsClient NewAttrStore func(string) AttrStore Logger logger.Logger - Txsrc string RowcacheOn bool + StorageConfig *storage.Config RBFConfig *rbfcfg.Config AntiEntropyInterval time.Duration } @@ -233,7 +233,7 @@ func DefaultHolderConfig() *HolderConfig { StatsClient: stats.NopStatsClient, NewAttrStore: newNopAttrStore, Logger: logger.NopLogger, - Txsrc: DefaultTxsrc, + StorageConfig: storage.NewDefaultConfig(), RBFConfig: rbfcfg.NewDefaultConfig(), } } @@ -247,9 +247,13 @@ func NewHolder(path string, cfg *HolderConfig) *Holder { if txsrc != "" { _ = MustTxsrcToTxtype(txsrc) // INVAR: have valid txsrc. - cfg.Txsrc = txsrc + cfg.StorageConfig.Backend = txsrc } - } else if cfg.RBFConfig == nil { + } + if cfg.StorageConfig == nil { + cfg.StorageConfig = storage.NewDefaultConfig() + } + if cfg.RBFConfig == nil { cfg.RBFConfig = rbfcfg.NewDefaultConfig() } @@ -271,7 +275,7 @@ func NewHolder(path string, cfg *HolderConfig) *Holder { OpenIDAllocator: cfg.OpenIDAllocator, translationSyncer: cfg.TranslationSyncer, Logger: cfg.Logger, - Opts: HolderOpts{Txsrc: cfg.Txsrc, RowcacheOn: cfg.RowcacheOn}, + Opts: HolderOpts{Txsrc: cfg.StorageConfig.Backend, RowcacheOn: cfg.RowcacheOn}, SnapshotQueue: defaultSnapshotQueue, @@ -282,9 +286,9 @@ func NewHolder(path string, cfg *HolderConfig) *Holder { indexes: make(map[string]*Index), } - rbf.SetRowcacheOn(cfg.RowcacheOn) + storage.SetRowCacheOn(cfg.RowcacheOn) - txf, err := NewTxFactory(cfg.Txsrc, path, h) + txf, err := NewTxFactory(cfg.StorageConfig.Backend, path, h) panicOn(err) h.txf = txf h.txf.blueGreenOffIfRunningBlueGreen() @@ -579,7 +583,7 @@ func (h *Holder) Open() error { defer func() { h.opening = false }() if h.txf == nil { - txf, err := NewTxFactory(h.cfg.Txsrc, h.path, h) + txf, err := NewTxFactory(h.cfg.StorageConfig.Backend, h.path, h) if err != nil { return errors.Wrap(err, "Holder.Open NewTxFactory()") } diff --git a/rbf.go b/rbf.go index dff5a1123..86af392a2 100644 --- a/rbf.go +++ b/rbf.go @@ -30,6 +30,8 @@ import ( rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg" "github.com/pilosa/pilosa/v2/roaring" txkey "github.com/pilosa/pilosa/v2/short_txkey" + "github.com/pilosa/pilosa/v2/storage" + //txkey "github.com/pilosa/pilosa/v2/txkey" "github.com/pkg/errors" ) @@ -90,6 +92,14 @@ type rbfDBRegistrar struct { mp map[*RbfDBWrapper]bool path2db map[string]*RbfDBWrapper + + rbfConfig *rbfcfg.Config +} + +func (r *rbfDBRegistrar) SetRBFConfig(cfg *rbfcfg.Config) { + r.mu.Lock() + defer r.mu.Unlock() + r.rbfConfig = cfg } func (r *rbfDBRegistrar) Size() int { @@ -148,7 +158,7 @@ func rbfPath(path string) string { // if one does not exist for its path. Otherwise it returns // the existing instance. This insures only one RbfDBWrapper // per bpath in this pilosa node. -func (r *rbfDBRegistrar) OpenDBWrapper(path0 string, doAllocZero bool, cfg *rbfcfg.Config) (DBWrapper, error) { +func (r *rbfDBRegistrar) OpenDBWrapper(path0 string, doAllocZero bool, cfg *storage.Config) (DBWrapper, error) { path := rbfPath(path0) r.mu.Lock() defer r.mu.Unlock() @@ -157,11 +167,12 @@ func (r *rbfDBRegistrar) OpenDBWrapper(path0 string, doAllocZero bool, cfg *rbfc // creates the effect of having only one DB open per pilosa node. return w, nil } - if cfg == nil { - cfg = rbfcfg.NewDefaultConfig() - cfg.DoAllocZero = doAllocZero + if r.rbfConfig == nil { + r.rbfConfig = rbfcfg.NewDefaultConfig() + r.rbfConfig.DoAllocZero = doAllocZero + r.rbfConfig.FsyncEnabled = cfg.FsyncEnabled } - db := rbf.NewDB(path, cfg) + db := rbf.NewDB(path, r.rbfConfig) w = &RbfDBWrapper{ reg: r, @@ -169,7 +180,7 @@ func (r *rbfDBRegistrar) OpenDBWrapper(path0 string, doAllocZero bool, cfg *rbfc db: db, doAllocZero: doAllocZero, openTx: make(map[*RBFTx]bool), - cfg: cfg, + cfg: r.rbfConfig, } r.unprotectedRegister(w) @@ -424,7 +435,7 @@ func (tx *RBFTx) UseRowCache() bool { // the rowCache without first making a copy. // So we only use the rowCache if the copy is // enabled. - return rbf.EnableRowCache() + return storage.EnableRowCache() } func (tx *RBFTx) ApplyFilter(index, field, view string, shard uint64, ckey uint64, filter roaring.BitmapFilter) (err error) { diff --git a/rbf/cfg/cfg.go b/rbf/cfg/cfg.go index f90179349..5f8cd9345 100644 --- a/rbf/cfg/cfg.go +++ b/rbf/cfg/cfg.go @@ -28,26 +28,26 @@ const ( type Config struct { // The maximum allowed database size. Required by mmap. - MaxSize int64 + MaxSize int64 `toml:"max-db-size"` // The maximum allowed WAL size. Required by mmap. - MaxWALSize int64 + MaxWALSize int64 `toml:"max-wal-size"` // The minimum WAL size before the WAL is copied to the DB. - MinWALCheckpointSize int64 + MinWALCheckpointSize int64 `toml:"min-wal-checkpoint-size"` // The maximum WAL size before transactions are halted to allow a checkpoint. - MaxWALCheckpointSize int64 + MaxWALCheckpointSize int64 `toml:"max-wal-checkpoint-size"` // Set before calling db.Open() - FsyncEnabled bool + FsyncEnabled bool `toml:"fsync"` // for mmap correctness testing. - DoAllocZero bool + DoAllocZero bool `toml:"do-alloc-zero"` // CursorCacheSize is the number of copies of Cursor{} to keep in our // readyCursorCh arena to avoid GC pressure. - CursorCacheSize int64 + CursorCacheSize int64 `toml:"cursor-cache-size"` } func NewDefaultConfig() *Config { @@ -66,13 +66,12 @@ func NewDefaultConfig() *Config { func (cfg *Config) DefineFlags(flags *pflag.FlagSet) { default0 := NewDefaultConfig() - flags.Int64Var(&cfg.MaxSize, "rbf-max-db-size", default0.MaxSize, "RBF maximum size in bytes of a database file (distinct from a WAL file)") - flags.Int64Var(&cfg.MaxWALSize, "rbf-max-wal-size", default0.MaxWALSize, "RBF maximum size in bytes of a WAL file (distinct from a DB file)") - flags.Int64Var(&cfg.MinWALCheckpointSize, "rbf-min-wal-checkpoint-size", default0.MinWALCheckpointSize, "RBF minimum size in bytes of a WAL file before attempting checkpoint") - flags.Int64Var(&cfg.MaxWALCheckpointSize, "rbf-max-wal-checkpoint-size", default0.MaxWALCheckpointSize, "RBF maximum size in bytes of a WAL file before forcing checkpoint") + flags.Int64Var(&cfg.MaxSize, "rbf.max-db-size", default0.MaxSize, "RBF maximum size in bytes of a database file (distinct from a WAL file)") + flags.Int64Var(&cfg.MaxWALSize, "rbf.max-wal-size", default0.MaxWALSize, "RBF maximum size in bytes of a WAL file (distinct from a DB file)") + flags.Int64Var(&cfg.MinWALCheckpointSize, "rbf.min-wal-checkpoint-size", default0.MinWALCheckpointSize, "RBF minimum size in bytes of a WAL file before attempting checkpoint") + flags.Int64Var(&cfg.MaxWALCheckpointSize, "rbf.max-wal-checkpoint-size", default0.MaxWALCheckpointSize, "RBF maximum size in bytes of a WAL file before forcing checkpoint") // renamed from --rbf-fsync to just --fsync because now it applies to all Tx backends. flags.BoolVar(&cfg.FsyncEnabled, "fsync", default0.FsyncEnabled, "enable fsync fully safe flush-to-disk") - flags.Int64Var(&cfg.CursorCacheSize, "rbf-cursor-cache", default0.CursorCacheSize, "how big a Cursor arena to maintain. 0 means use sync.Pool with dynamic sizing. Note that <= 20 is needed to pass CI. Controls the memory footprint of rbf.") - + flags.Int64Var(&cfg.CursorCacheSize, "rbf.cursor-cache-size", default0.CursorCacheSize, "how big a Cursor arena to maintain. 0 means use sync.Pool with dynamic sizing. Note that <= 20 is needed to pass CI. Controls the memory footprint of rbf.") } diff --git a/rbf/cursorx.go b/rbf/cursorx.go index 171bbd376..0ca291ff1 100644 --- a/rbf/cursorx.go +++ b/rbf/cursorx.go @@ -19,31 +19,13 @@ import ( "io" "math" "os" - "sync/atomic" "unsafe" "github.com/pilosa/pilosa/v2/roaring" + "github.com/pilosa/pilosa/v2/storage" "github.com/pkg/errors" ) -// if enableRowCache, then we must not return mmap-ed memory -// directly, but only a copy. -var enableRowcache int64 = 1 - -// SetEnableRowCache should only be called in NewHolder before -// all other reads. -func SetRowcacheOn(on bool) { - if on { - atomic.StoreInt64(&enableRowcache, 1) - } else { - atomic.StoreInt64(&enableRowcache, 0) - } -} - -func EnableRowCache() bool { - return atomic.LoadInt64(&enableRowcache) == 1 -} - //probably should just implement the container interface // but for now i'll do it func (c *Cursor) Rows() ([]uint64, error) { @@ -192,7 +174,7 @@ func intoContainer(l leafCell, tx *Tx, replacing *roaring.Container, target []by orig := l.Data var cpMaybe []byte var mapped bool - if EnableRowCache() || tx.db.cfg.DoAllocZero { + if storage.EnableRowCache() || tx.db.cfg.DoAllocZero { // make a copy, otherwise the rowCache will see corrupted data // or mmapped data that may disappear. cpMaybe = target[:len(orig)] @@ -209,7 +191,7 @@ func intoContainer(l leafCell, tx *Tx, replacing *roaring.Container, target []by case ContainerTypeBitmapPtr: _, bm, _ := tx.leafCellBitmap(toPgno(cpMaybe)) cloneMaybe := bm - if EnableRowCache() { + if storage.EnableRowCache() { cloneMaybe = (*[1024]uint64)(unsafe.Pointer(&target[0]))[:1024] copy(cloneMaybe, bm) } @@ -235,7 +217,7 @@ func toContainer(l leafCell, tx *Tx) (c *roaring.Container) { orig := l.Data var cpMaybe []byte var mapped bool - if EnableRowCache() || tx.db.cfg.DoAllocZero { + if storage.EnableRowCache() || tx.db.cfg.DoAllocZero { // make a copy, otherwise the rowCache will see corrupted data // or mmapped data that may disappear. cpMaybe = make([]byte, len(orig)) @@ -252,7 +234,7 @@ func toContainer(l leafCell, tx *Tx) (c *roaring.Container) { case ContainerTypeBitmapPtr: _, bm, _ := tx.leafCellBitmap(toPgno(cpMaybe)) cloneMaybe := bm - if EnableRowCache() { + if storage.EnableRowCache() { cloneMaybe = make([]uint64, len(bm)) copy(cloneMaybe, bm) } diff --git a/rrtx.go b/rrtx.go index 8c75a2924..fbe78eacf 100644 --- a/rrtx.go +++ b/rrtx.go @@ -26,10 +26,9 @@ import ( "sync" "sync/atomic" - "github.com/pilosa/pilosa/v2/rbf" - rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg" "github.com/pilosa/pilosa/v2/roaring" txkey "github.com/pilosa/pilosa/v2/short_txkey" + "github.com/pilosa/pilosa/v2/storage" //txkey "github.com/pilosa/pilosa/v2/txkey" "github.com/pkg/errors" @@ -68,7 +67,7 @@ func (tx *RoaringTx) Dump(short bool, shard uint64) { } func (tx *RoaringTx) UseRowCache() bool { - return rbf.EnableRowCache() + return storage.EnableRowCache() } // based on view.openFragments() @@ -637,8 +636,7 @@ func (r *roaringRegistrar) unregister(w *RoaringWrapper) { // openRoaringDB will check the registry and make a new instance only // if one does not exist for its path0. Otherwise it returns // the existing instance. -func (r *roaringRegistrar) OpenDBWrapper(path string, doAllocZero bool, cfg *rbfcfg.Config) (DBWrapper, error) { - +func (r *roaringRegistrar) OpenDBWrapper(path string, doAllocZero bool, _ *storage.Config) (DBWrapper, error) { r.mu.Lock() defer r.mu.Unlock() w, ok := r.path2db[path] diff --git a/server.go b/server.go index 5a775d61d..746c37843 100644 --- a/server.go +++ b/server.go @@ -36,6 +36,7 @@ import ( rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg" "github.com/pilosa/pilosa/v2/roaring" "github.com/pilosa/pilosa/v2/stats" + "github.com/pilosa/pilosa/v2/storage" "github.com/pilosa/pilosa/v2/topology" "github.com/pkg/errors" "golang.org/x/sync/errgroup" @@ -360,13 +361,12 @@ func OptServerOpenTranslateReader(fn OpenTranslateReaderFunc) ServerOption { } } -// OptServerTxsrc is a functional option on Server -// used to specify the transactional-storage to use, -// resulting in RoaringTx, RbfTx, BadgerTx, or a blueGreen* Tx -// being used for all Tx interface calls. -func OptServerTxsrc(txsrc string) ServerOption { +// OptServerStorageConfig is a functional option on Server used to specify the +// transactional-storage backend to use, resulting in RoaringTx, RbfTx, +// BadgerTx, or a blueGreen* Tx being used for all Tx interface calls. +func OptServerStorageConfig(cfg *storage.Config) ServerOption { return func(s *Server) error { - s.holderConfig.Txsrc = txsrc + s.holderConfig.StorageConfig = cfg return nil } } diff --git a/server/config.go b/server/config.go index 900d6ddeb..920320d97 100644 --- a/server/config.go +++ b/server/config.go @@ -27,6 +27,7 @@ import ( petcd "github.com/pilosa/pilosa/v2/etcd" "github.com/pilosa/pilosa/v2/gossip" rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg" + "github.com/pilosa/pilosa/v2/storage" "github.com/pilosa/pilosa/v2/toml" "github.com/pkg/errors" ) @@ -206,18 +207,20 @@ type Config struct { // returned from the blueGreenTx. Txsrc string `toml:"txsrc"` + Storage *storage.Config `toml:"storage"` + // RowcacheOn, if true, turns on the row cache for all storage backends. // The default is now off because it makes rbf queries faster and uses // much less memory. RowcacheOn bool `toml:"rowcache-on"` // RBFConfig defines all externally configurable RBF flags. - RBFConfig *rbfcfg.Config + RBFConfig *rbfcfg.Config `toml:"rbf"` // QueryHistoryLength sets the maximum number of queries that are maintained // for the /query-history endpoint. This parameter is per-node, and the // result combines the history from all nodes. - QueryHistoryLength int + QueryHistoryLength int `toml:"query-history-length"` } // MustValidate checks that all ports in a Config are unique and not zero. @@ -308,6 +311,7 @@ func NewConfig() *Config { WorkerPoolSize: runtime.NumCPU(), ImportWorkerPoolSize: runtime.NumCPU(), + Storage: storage.NewDefaultConfig(), RBFConfig: rbfcfg.NewDefaultConfig(), QueryHistoryLength: 100, diff --git a/server/server.go b/server/server.go index c3613ed0d..0709aeddf 100644 --- a/server/server.go +++ b/server/server.go @@ -441,7 +441,7 @@ func (m *Command) SetupServer() error { pilosa.OptServerClusterDisabled(m.Config.Cluster.Disabled, m.Config.Cluster.Hosts), pilosa.OptServerClusterName(m.Config.Cluster.Name), pilosa.OptServerSerializer(proto.Serializer{}), - pilosa.OptServerTxsrc(m.Config.Txsrc), + pilosa.OptServerStorageConfig(m.Config.Storage), pilosa.OptServerRowcacheOn(m.Config.RowcacheOn), pilosa.OptServerRBFConfig(m.Config.RBFConfig), pilosa.OptServerQueryHistoryLength(m.Config.QueryHistoryLength), diff --git a/storage/cache.go b/storage/cache.go new file mode 100644 index 000000000..7fda5af83 --- /dev/null +++ b/storage/cache.go @@ -0,0 +1,36 @@ +// Copyright 2017 Pilosa Corp. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +package storage + +import ( + "sync/atomic" +) + +// if enableRowCache, then we must not return mmap-ed memory +// directly, but only a copy. +var enableRowcache int64 = 1 + +// SetRowCacheOn should only be called in NewHolder before +// all other reads. +func SetRowCacheOn(on bool) { + if on { + atomic.StoreInt64(&enableRowcache, 1) + } else { + atomic.StoreInt64(&enableRowcache, 0) + } +} + +func EnableRowCache() bool { + return atomic.LoadInt64(&enableRowcache) == 1 +} diff --git a/storage/config.go b/storage/config.go new file mode 100644 index 000000000..68137a878 --- /dev/null +++ b/storage/config.go @@ -0,0 +1,42 @@ +// Copyright 2017 Pilosa Corp. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package storage + +// public strings that pilosa/server/config.go can reference +const ( + RoaringBackend string = "roaring" + RBFBackend string = "rbf" + BoltBackend string = "bolt" +) + +// DefaultBackend is set here. pilosa/server/config.go references it +// to set the default for pilosa server exeutable. +const DefaultBackend = RoaringBackend + +// Config represents configuration which applies to multiple storage engines. +type Config struct { + Backend string `toml:"backend"` + + // Set before calling db.Open() + FsyncEnabled bool `toml:"fsync"` +} + +// NewDefaultConfig returns a new Config with default values. +func NewDefaultConfig() *Config { + return &Config{ + Backend: DefaultBackend, + FsyncEnabled: true, + } +} diff --git a/test/pilosa.go b/test/pilosa.go index e999a3cd3..b5335ca23 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -70,6 +70,12 @@ func newCommand(tb testing.TB, opts ...server.CommandOption) *Command { m.Config.DataDir = path defaultConf := server.NewConfig() + // TODO: this is temporary and should be removed and + // automatically replaced with PILOSA_STORAGE_BACKEND. + if txsrc := os.Getenv("PILOSA_TXSRC"); txsrc != "" { + m.Config.Storage.Backend = txsrc + } + if m.Config.Bind == defaultConf.Bind { m.Config.Bind = "http://localhost:0" } diff --git a/txfactory.go b/txfactory.go index d8aaf3aee..92018c03f 100644 --- a/txfactory.go +++ b/txfactory.go @@ -28,9 +28,9 @@ import ( "text/tabwriter" "github.com/pilosa/pilosa/v2/hash" - "github.com/pilosa/pilosa/v2/rbf" "github.com/pilosa/pilosa/v2/roaring" txkey "github.com/pilosa/pilosa/v2/short_txkey" + "github.com/pilosa/pilosa/v2/storage" //txkey "github.com/pilosa/pilosa/v2/txkey" "github.com/pkg/errors" "github.com/zeebo/blake3" @@ -521,7 +521,7 @@ func NewTxFactory(txsrc string, holderDir string, holder *Holder) (f *TxFactory, f.blueGreenReg = newBlueGreenReg(types) f.isBlueGreen = true // blue-green can never use the rowCache. - rbf.SetRowcacheOn(false) + storage.SetRowCacheOn(false) } f.dbPerShard = f.NewDBPerShard(types, holderDir, holder) @@ -544,7 +544,7 @@ func (f *TxFactory) Open() error { // to determine if it should use the rowCache. Currently it // doesn't have a tx Tx parameter, so we use the Txf instead. func (f *TxFactory) UseRowCache() bool { - return rbf.EnableRowCache() + return storage.EnableRowCache() } // Txo holds the transaction options