mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-05 16:15:56 +00:00
- enable row cache again. Was off for tx perf measurement. - centralize UseRowCache choice to just one point, in rbf.EnableRowCache
1924 lines
58 KiB
Go
1924 lines
58 KiB
Go
// Copyright 2020 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 pilosa
|
|
|
|
import (
|
|
"bytes"
|
|
"fmt"
|
|
"io"
|
|
"io/ioutil"
|
|
"log"
|
|
"math"
|
|
"os"
|
|
"sort"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
badger "github.com/dgraph-io/badger/v2"
|
|
badgeroptions "github.com/dgraph-io/badger/v2/options"
|
|
"github.com/pilosa/pilosa/v2/hash"
|
|
"github.com/pilosa/pilosa/v2/rbf"
|
|
"github.com/pilosa/pilosa/v2/roaring"
|
|
"github.com/pilosa/pilosa/v2/testhook"
|
|
"github.com/pilosa/pilosa/v2/txkey"
|
|
"github.com/pkg/errors"
|
|
)
|
|
|
|
// TODO: is there a more optimal time to do badger garbage collection?
|
|
// As in: do we need to be more aggressive about cleaning in
|
|
// proportion to write activity? Space monitoring available with the
|
|
// badger.DB.Size() (lsm, vlog int64) call.
|
|
//
|
|
// See: https://godoc.org/github.com/dgraph-io/badger#DB.RunValueLogGC
|
|
// and: https://github.com/dgraph-io/badger#garbage-collection
|
|
//
|
|
// For now we run GC periodically every 1 minute or as set by the
|
|
// BadgerDBWrapper.GcEveryDur duration.
|
|
//
|
|
// Background: (quoting from docs referenced above)
|
|
//
|
|
// "Badger relies on the client to perform garbage collection at a time of
|
|
// their choosing. It provides the following method, which can be invoked
|
|
// at an appropriate time:
|
|
//
|
|
// "DB.RunValueLogGC(): This method is designed to do garbage collection while
|
|
// Badger is online. Along with randomly picking a file, it uses statistics
|
|
// generated by the LSM-tree compactions to pick files that are likely to
|
|
// lead to maximum space reclamation. It is recommended to be called during
|
|
// periods of low activity in your system, or periodically. One call would
|
|
// only result in removal of at max one log file. As an optimization, you
|
|
// could also immediately re-run it whenever it returns nil error (indicating
|
|
// a successful value log GC), as shown below."
|
|
//
|
|
// ticker := time.NewTicker(5 * time.Minute)
|
|
// defer ticker.Stop()
|
|
// for range ticker.C {
|
|
// again:
|
|
// err := db.RunValueLogGC(0.5)
|
|
// if err == nil {
|
|
// goto again
|
|
// }
|
|
// }
|
|
//
|
|
|
|
//
|
|
// =========================================================
|
|
// A note on using a recent version of badgerdb:
|
|
//
|
|
// We require a v2 release of badger after 2020 May 13, when support for
|
|
// multiple read-write iterators within one transaction was added.
|
|
// Many executor_test.go tests do foreachRow() operations,
|
|
// which call BadgerTx.ContainerIterator(), which in turn creates
|
|
// a first read-write iterator, and then OffsetRange(), which needs a
|
|
// second iterator, while still in the same read-write transaction.
|
|
//
|
|
// The most recent v2 master was pulled in and added to go.mod
|
|
// by doing go get github.com/dgraph-io/badger/v2@master
|
|
// resulting in the go.mod line
|
|
// github.com/dgraph-io/badger/v2 v2.0.1-rc1.0.20200709123515-8e896a7af361
|
|
// as of this writing, 2020 July 09. This version contains the support
|
|
// for having multiple read-write iterators.
|
|
//
|
|
// Reference on github.com/dgraph-io/badger
|
|
//
|
|
// commit af22dfd8d51317d765f0c05dcdf1d15981cca4f3
|
|
// Author: Elliot Courant <me@elliotcourant.dev>
|
|
// Date: Wed May 13 01:07:33 2020 -0500
|
|
//
|
|
// Support multiple iterators in read-write transactions. (#1286)
|
|
//
|
|
// This adds support for multiple iterators during a read-write transaction. The
|
|
// iterators created in a read-write transaction will only be able to see writes
|
|
// that were performed before the iterator was created. Any writes that occur
|
|
// after the iterator is created will be invisible to the iterator.
|
|
//
|
|
// Fixes https://github.com/dgraph-io/badger/issues/981
|
|
//
|
|
//
|
|
// Otherwise we'll get these panics:
|
|
// 'Only one iterator can be active at one time, for a RW txn.'
|
|
// when trying to open a second iterator on the same write transaction.
|
|
// e.g. go test -v -run TestExecutor_TranslateRowsOnBool
|
|
|
|
var badgerDefaultLogger *BadgerLog
|
|
var badgerTestLogger *BadgerLog
|
|
|
|
const BadgerLogToStderr = false
|
|
|
|
func init() {
|
|
// badger test output clutters up the screen, dump to /dev/null for now.
|
|
null, err := os.Open(os.DevNull)
|
|
panicOn(err)
|
|
var out io.Writer = null
|
|
if BadgerLogToStderr {
|
|
// view badger logs
|
|
out = os.Stderr
|
|
}
|
|
badgerTestLogger = &BadgerLog{Logger: log.New(out, "badger ", log.LstdFlags)}
|
|
badgerDefaultLogger = badgerTestLogger
|
|
|
|
// BadgerDB recommends a minimum of 128 GOMAXPROCS to make use of the IOPs
|
|
// available on the SSD. So we set that here. Details:
|
|
//
|
|
// from https://github.com/dgraph-io/badger#are-there-any-go-specific-settings-that-i-should-use
|
|
//
|
|
// "We *highly* recommend setting a high number for GOMAXPROCS,
|
|
// which allows Go to observe the full IOPS throughput provided by
|
|
// modern SSDs. In Dgraph, we have set it to 128. For more details,
|
|
// see this thread [https://groups.google.com/forum/#!topic/golang-nuts/jPb_h3TvlKE/discussion]."
|
|
//
|
|
// From that thread on golang-nuts:
|
|
//
|
|
// "Manish Rai Jain
|
|
// 8/7/17
|
|
// Hey folks,
|
|
// During Gophercon, I happened to meet Russ Cox and asked him the same question.
|
|
// If File::Read blocks goroutines, which then spawn new OS threads, in a long running job,
|
|
// there should be plenty of OS threads created already, so the random read throughput
|
|
// should increase over time and stabilize to the maximum possible value. But, that's
|
|
// not what I see in my benchmarks.
|
|
//
|
|
// And his explanation was that the GOMAXPROCS in a way acts like a multiplexer.
|
|
// From docs, "the GOMAXPROCS variable limits the number of operating system threads
|
|
// that can execute user-level Go code simultaneously." Which basically means, all
|
|
// reads must first be run only via GOMAXPROCS number of goroutines, before switching
|
|
// over to some OS thread (not really a switch, but conceptually speaking). This
|
|
// introduces a bottleneck for throughput.
|
|
// I re-ran my benchmarks with a much higher GOMAXPROCS and was able to then
|
|
// achieve the maximum throughput. The numbers are here:
|
|
// https://github.com/dgraph-io/badger-bench/blob/master/randread/maxprocs.txt
|
|
// To summarize these benchmarks, Linux fio achieves 118K IOPS, and with GOMAXPROCS=64/128,
|
|
// I'm able to achieve 105K IOPS, which is close enough. Win!
|
|
//
|
|
// Regarding the point about using io_submit etc., instead of goroutines; I managed to
|
|
// find a library which does that, but it performed worse than just using goroutines.
|
|
// https://github.com/traetox/goaio/issues/3
|
|
// From what I gather (talking to Russ and Ian), whatever work is going on in user space,
|
|
// the same work has to happen in kernel space; so there's not much benefit here.
|
|
//
|
|
// Overall, with GOMAXPROCS set to a higher value (as I've done in Dgraph), one can get
|
|
// the advertised SSD throughput using goroutines."
|
|
//
|
|
// moved to holder init
|
|
//runtime.GOMAXPROCS(128)
|
|
}
|
|
|
|
// BadgerLog exists because badger requires a particular logger interface, with a
|
|
// Debugf method that is not on standard library log.Logger
|
|
type BadgerLog struct {
|
|
*log.Logger
|
|
}
|
|
|
|
// Errorf logs an error.
|
|
func (l *BadgerLog) Errorf(f string, v ...interface{}) {
|
|
l.Printf("ERROR: "+f, v...)
|
|
}
|
|
|
|
// Warningf logs a warning.
|
|
func (l *BadgerLog) Warningf(f string, v ...interface{}) {
|
|
l.Printf("WARNING: "+f, v...)
|
|
}
|
|
|
|
// Infof logs an informational statement.
|
|
func (l *BadgerLog) Infof(f string, v ...interface{}) {
|
|
l.Printf("INFO: "+f, v...)
|
|
}
|
|
|
|
// Debugf logs a debug statement.
|
|
func (l *BadgerLog) Debugf(f string, v ...interface{}) {
|
|
l.Printf("DEBUG: "+f, v...)
|
|
}
|
|
|
|
// badgerRegistrar facilitates shutdown
|
|
// of all the badger databases started under
|
|
// tests. Its needed because most tests don't cleanup
|
|
// the *Index(es) they create. But we still
|
|
// want to shutdown badgerDB goroutines
|
|
// after tests run.
|
|
//
|
|
// It also allows opening the same path twice to
|
|
// result in sharing the same open database handle, and
|
|
// thus the same transactional guarantees.
|
|
//
|
|
type badgerRegistrar struct {
|
|
mu sync.Mutex
|
|
mp map[*BadgerDBWrapper]bool
|
|
|
|
path2db map[string]*BadgerDBWrapper
|
|
}
|
|
|
|
func (r *badgerRegistrar) Size() int {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
nmp := len(r.mp)
|
|
npa := len(r.path2db)
|
|
if nmp != npa {
|
|
panic(fmt.Sprintf("nmp=%v, vs npa=%v", nmp, npa))
|
|
}
|
|
return nmp
|
|
}
|
|
|
|
var globalBadgerReg *badgerRegistrar = newBadgerTestRegistrar()
|
|
|
|
func newBadgerTestRegistrar() *badgerRegistrar {
|
|
return &badgerRegistrar{
|
|
mp: make(map[*BadgerDBWrapper]bool),
|
|
path2db: make(map[string]*BadgerDBWrapper),
|
|
}
|
|
}
|
|
|
|
// register each badger created under tests, so we
|
|
// can clean them up. This is called by openBadgerDBWrapper() while
|
|
// holding the r.mu.Lock, since it needs to atomically
|
|
// check the registry and make a new instance only
|
|
// if one does not exist for its path, and otherwise
|
|
// return the existing instance.
|
|
func (r *badgerRegistrar) unprotectedRegister(w *BadgerDBWrapper) {
|
|
r.mp[w] = true
|
|
r.path2db[w.path] = w
|
|
}
|
|
|
|
// unregister removes w from r
|
|
func (r *badgerRegistrar) unregister(w *BadgerDBWrapper) {
|
|
r.mu.Lock()
|
|
delete(r.mp, w)
|
|
delete(r.path2db, w.path)
|
|
r.mu.Unlock()
|
|
}
|
|
|
|
func DumpAllBadger() {
|
|
globalBadgerReg.mu.Lock()
|
|
defer globalBadgerReg.mu.Unlock()
|
|
for w := range globalBadgerReg.mp {
|
|
_ = w
|
|
AlwaysPrintf("this badger path='%v' has: \n%v\n", w.path, w.StringifiedBadgerKeys(nil, true))
|
|
}
|
|
}
|
|
|
|
// badgerPath is a helper for determining the full directory
|
|
// in which the badger database will be stored.
|
|
func badgerPath(path string) string {
|
|
if !strings.HasSuffix(path, "-badgerdb@") {
|
|
return path + "-badgerdb@"
|
|
}
|
|
return path
|
|
}
|
|
|
|
// openBadgerDB opens the database in the bpath directoy
|
|
// without deleting any prior content. Any BadgerDB
|
|
// database directory will have the "-badgerdb" suffix.
|
|
//
|
|
// openBadgerDB will check the registry and make a new instance only
|
|
// if one does not exist for its bpath. Otherwise it returns
|
|
// the existing instance. This insures only one badgerDB
|
|
// per bpath in this pilosa node.
|
|
func (r *badgerRegistrar) OpenDBWrapper(bpath string, doAllocZero bool) (DBWrapper, error) {
|
|
// now that newTxFactory can call us directly, we might not
|
|
// have the -badgerdb suffix.
|
|
bpath = badgerPath(bpath)
|
|
|
|
err := os.MkdirAll(bpath, 0755)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
w, ok := r.path2db[bpath]
|
|
if ok {
|
|
// creates the effect of having only one badger open per pilosa node.
|
|
return w, nil
|
|
}
|
|
// otherwise, make a new badger and store it in globalBadgerReg
|
|
|
|
// regular: works on amd64, but 386 doesn't work.
|
|
opt := badger.DefaultOptions(bpath).WithLogger(badgerDefaultLogger)
|
|
|
|
opt.Compression = badgeroptions.None // turn off compression.
|
|
opt.ZSTDCompressionLevel = 0 // really, just in case.
|
|
opt.SyncWrites = true // default is true, safe.
|
|
//opt.KeepL0InMemory = true // speedup?
|
|
|
|
// MaxCacheSize docs:
|
|
//
|
|
// how much data cache should hold in memory. A small size of
|
|
// cache means lower memory consumption and lookups/iterations
|
|
// would take longer. It is recommended to use a cache if you're
|
|
// using compression or encryption. If compression and
|
|
// encryption both are disabled, adding a cache will lead to
|
|
// unnecessary overhead which will affect the read performance.
|
|
// Setting size to zero disables the cache altogether.
|
|
//opt.MaxCacheSize = 1 << 30 // slows down 135 sec vs 113 sec on our benchmark
|
|
//opt.MaxCacheSize = 0
|
|
opt.LoadBloomsOnOpen = false // should speed up start-up time.
|
|
|
|
//opt.KeepBlocksInCache = true // default false
|
|
//opt.KeepBlockIndicesInCache = true // default false
|
|
|
|
opt.BlockSize = 8 * 1024 // default 4 * 1024
|
|
|
|
// to get memory only do:
|
|
//opt := badger.DefaultOptions("").WithLogger(badgerDefaultLogger).WithInMemory(true)
|
|
|
|
//db, err := badger.Open(opt)
|
|
db, err := badger.OpenManaged(opt)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
halt := make(chan bool)
|
|
w = &BadgerDBWrapper{
|
|
reg: r,
|
|
path: bpath,
|
|
db: db,
|
|
halt: halt,
|
|
hasher: hash.NewBlake3Hasher(),
|
|
doAllocZero: doAllocZero,
|
|
}
|
|
_ = testhook.Opened(NewAuditor(), w, nil)
|
|
r.unprotectedRegister(w)
|
|
|
|
w.startBadgerGarbageCollectionBackgroundGoro()
|
|
return w, nil
|
|
}
|
|
|
|
func (w *BadgerDBWrapper) DeleteDBPath(dbs *DBShard) error {
|
|
path := dbs.pathForType(badgerTxn)
|
|
return os.RemoveAll(path)
|
|
}
|
|
|
|
// DeleteIndex deletes all the containers associated with
|
|
// the named index from the badger database.
|
|
func (w *BadgerDBWrapper) DeleteIndex(indexName string) error {
|
|
|
|
// We use the apostrophie rune `'` to locate the end of the
|
|
// index name in the key prefix, so we cannot allow indexNames
|
|
// themselves to contain apostrophies.
|
|
if strings.Contains(indexName, "'") {
|
|
return fmt.Errorf("error: bad indexName `%v` in BadgerDBWrapper.DeleteIndex() call: indexName cannot contain apostrophes/single quotes.", indexName)
|
|
}
|
|
prefix := txkey.IndexOnlyPrefix(indexName)
|
|
return w.DeletePrefix(prefix)
|
|
}
|
|
|
|
// startBadgerGarbageCollectionBackgroundGoro handles Badger DB
|
|
// garbage colection by regularly purging the value log from
|
|
// a background goroutine. w.GcEveryDur controls how often
|
|
// it runs. The default is after every 60 seconds.
|
|
func (w *BadgerDBWrapper) startBadgerGarbageCollectionBackgroundGoro() {
|
|
go func() {
|
|
dur := w.GcEveryDur
|
|
if dur == 0 {
|
|
dur = time.Minute
|
|
}
|
|
ticker := time.NewTicker(dur)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ticker.C:
|
|
w.muGC.Lock()
|
|
again:
|
|
err := w.db.RunValueLogGC(0.5)
|
|
if err == nil {
|
|
goto again
|
|
}
|
|
w.muGC.Unlock()
|
|
case <-w.halt:
|
|
return
|
|
}
|
|
}
|
|
}()
|
|
}
|
|
|
|
// statically confirm that BadgerTx satisfies the Tx interface.
|
|
var _ Tx = (*BadgerTx)(nil)
|
|
|
|
var globalNextTxSnBadger int64
|
|
|
|
// BadgerDBWrapper provides the NewTx() method.
|
|
// The methods on BadgerDBWrapper are thread-safe, and can be called
|
|
// from different goroutines/threads.
|
|
type BadgerDBWrapper struct {
|
|
// serialize operations on BadgerDBWrapper and thus on the .db too,
|
|
// when obtaining new txns on different goroutines.
|
|
muDb sync.Mutex
|
|
|
|
path string
|
|
db *badger.DB
|
|
|
|
// track our registrar for Close / goro leak reporting purposes.
|
|
reg *badgerRegistrar
|
|
|
|
// openTx and openIt are BadgerDBWrapper scoped tables of all open
|
|
// transactions and iterators. These are primarily for debugging purposes.
|
|
// openTx and openIt should only be read/written after locking the muOpenTxIt mutex.
|
|
|
|
// the bool value is the writable attribute of the key *BadgerTx
|
|
openTx map[*BadgerTx]bool
|
|
|
|
// the bool value is whether the iterator is reversed
|
|
openIt map[*BadgerIterator]bool
|
|
|
|
// protect openTx and openIt
|
|
muOpenTxIt sync.Mutex
|
|
|
|
// close(halt) to shutdown the badger gc goroutine in Close()
|
|
halt chan bool
|
|
|
|
// make BadgerDBWrapper.Close() idempotent, avoiding panic on double Close()
|
|
closed bool
|
|
|
|
// GcEveryDur controls how often the background goroutine
|
|
// runs garbage collection on the on-disk values-log.
|
|
// It defaults to running a GC every 1 minute if left as 0.
|
|
GcEveryDur time.Duration
|
|
|
|
// muGC ensures we only run one Garbage Collection at a time.
|
|
muGC sync.Mutex
|
|
|
|
hasher *hash.Blake3Hasher
|
|
|
|
// doAllocZero sets the corresponding flag on all new BadgerTx.
|
|
// When doAllocZero is true, we zero out any data from badger
|
|
// after transcation commit and rollback. This simulates
|
|
// what would happen if we were to use the mmap-ed data
|
|
// from badger directly. Currently we copy by default for
|
|
// safety because otherwise TestAPI_ImportColumnAttrs sees
|
|
// corrupted data.
|
|
doAllocZero bool
|
|
|
|
// stack() from our creation point, to track tests
|
|
// that haven't closed us.
|
|
startStack string
|
|
|
|
DeleteEmptyContainer bool
|
|
|
|
writeBatch *badger.WriteBatch
|
|
}
|
|
|
|
func (w *BadgerDBWrapper) Path() string {
|
|
return w.path
|
|
}
|
|
|
|
func (w *BadgerDBWrapper) SetHolder(h *Holder) {
|
|
// don't need it at the moment
|
|
//w.h = h
|
|
}
|
|
|
|
// EntryCount for Badger returns a count of either 0 or 1. A 1
|
|
// means it
|
|
func (w *BadgerDBWrapper) HasData() (has bool, err error) {
|
|
tx, err := w.NewTx(!writable, "", Txo{Write: false})
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
defer tx.Rollback()
|
|
bi := NewBadgerIterator(tx.(*BadgerTx), nil)
|
|
defer bi.Close()
|
|
bi.Seek(nil)
|
|
if bi.it.Valid() {
|
|
return true, nil
|
|
}
|
|
return false, nil
|
|
}
|
|
|
|
func (w *BadgerDBWrapper) CleanupTx(tx Tx) {
|
|
r := tx.(*BadgerTx)
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
w.unprotectedCleanupTx(r)
|
|
}
|
|
|
|
func (w *BadgerDBWrapper) unprotectedCleanupTx(btx *BadgerTx) {
|
|
if btx.done {
|
|
return
|
|
}
|
|
btx.done = true
|
|
btx.o.dbs.Cleanup(btx) // release the read/write lock.
|
|
}
|
|
|
|
// unprotectedListOpenTxAsString is a debugging helper.
|
|
// It is not thread safe, but is only used for debugging. Called internally while
|
|
// holding locks.
|
|
func (w *BadgerDBWrapper) unprotectedListOpenTxAsString() (r string) {
|
|
|
|
r = "openTx list = ["
|
|
for txn, write := range w.openTx {
|
|
r += fmt.Sprintf("txn p=%p(write:%v), ", txn, write)
|
|
}
|
|
return r + "]"
|
|
}
|
|
|
|
var _ = (*BadgerDBWrapper)(nil).unprotectedListOpenTxAsString // linter happy
|
|
|
|
// UnprotectedListOpenItAsString is exported because it is
|
|
// used for debugging in some of the pilosa_test tests.
|
|
// It is not thread safe, but only used for debugging. Called internally
|
|
// while holding locks and externally while not.
|
|
func (w *BadgerDBWrapper) UnprotectedListOpenItAsString() (r string) {
|
|
r = "openIt list = ["
|
|
for it, reverse := range w.openIt {
|
|
r += fmt.Sprintf("it p=%p(reverse:%v), ", it, reverse)
|
|
}
|
|
return r + "]"
|
|
}
|
|
|
|
// NewTx produces BadgerDB based ACID transactions. If
|
|
// the transaction will modify data, then the write flag must be true.
|
|
// Read-only queries should set write to false, to allow more concurrency.
|
|
// Methods on a BadgerTx are thread-safe, and can be called from
|
|
// different goroutines.
|
|
//
|
|
// initialIndexName is optional. It is set by the TxFactory from the Txo
|
|
// options provided at the Tx creation point. It allows us to recognize
|
|
// and isolate cross-index queries more quickly. It can always be empty ""
|
|
// but when set is highly useful for debugging. It has no impact
|
|
// on transaction behavior.
|
|
//
|
|
func (w *BadgerDBWrapper) NewTx(write bool, initialIndexName string, o Txo) (tx Tx, err error) {
|
|
|
|
sn := atomic.AddInt64(&globalNextTxSnBadger, 1)
|
|
|
|
tx = &BadgerTx{
|
|
frag: o.Fragment,
|
|
write: write,
|
|
tx: w.db.NewTransactionAt(1, write),
|
|
Db: w,
|
|
doAllocZero: w.doAllocZero,
|
|
initialIndexName: initialIndexName,
|
|
DeleteEmptyContainer: w.DeleteEmptyContainer,
|
|
o: o,
|
|
sn: sn,
|
|
}
|
|
return
|
|
}
|
|
|
|
// Close shuts down the Badger database.
|
|
func (w *BadgerDBWrapper) Close() (err error) {
|
|
|
|
w.muDb.Lock()
|
|
defer w.muDb.Unlock()
|
|
if !w.closed {
|
|
w.reg.unregister(w)
|
|
close(w.halt)
|
|
w.closed = true
|
|
}
|
|
_ = testhook.Closed(NewAuditor(), w, nil)
|
|
return w.db.Close()
|
|
}
|
|
|
|
// BadgerTx wraps a badger.Txn and provides the Tx interface
|
|
// method implementations.
|
|
// The methods on BadgerTx are thread-safe, and can be called
|
|
// from different goroutines.
|
|
type BadgerTx struct {
|
|
|
|
// mu serializes badger operations on this single txn instance.
|
|
//
|
|
// reference: https://godoc.org/github.com/dgraph-io/badger
|
|
// "Running [two separate -jea] transactions concurrently is OK. However, a
|
|
// transaction itself isn't thread safe, and should only
|
|
// be run serially. It doesn't matter if a transaction is
|
|
// created by one goroutine and passed down to other, as
|
|
// long as the Txn APIs are called serially."
|
|
mu sync.Mutex
|
|
|
|
write bool
|
|
Db *BadgerDBWrapper
|
|
tx *badger.Txn
|
|
|
|
frag *fragment
|
|
opcount int
|
|
|
|
doAllocZero bool
|
|
|
|
// for tracking txn boundary issues, track all the memory
|
|
// that we deploy for roaring containers, and zero it on
|
|
// transaction commit/rollback.
|
|
acMu sync.Mutex // protect ourAllocs and ourContainers
|
|
ourAllocs [][]byte
|
|
ourContainers []*roaring.Container
|
|
|
|
initialIndexName string
|
|
|
|
DeleteEmptyContainer bool
|
|
|
|
// We must avoid writing more than 10MB to badger in
|
|
// one transaction. If we go over, then
|
|
// we'll get a ErrTxnTooBig error. At that point
|
|
// we can't commit more, because the transaction
|
|
// will "conflict". So we must monitor
|
|
// totals written and auto-commit before going
|
|
// over the limits to avoid wedging into an
|
|
// unrecoverable state.
|
|
//
|
|
// We now use the approach described in https://github.com/dgraph-io/badger/issues/539
|
|
// wherein "I fixed it by using a ManagedDB and using a commtTs of 1 always."
|
|
// But see also
|
|
// https://github.com/dgraph-io/badger/issues/409
|
|
// wrt to "Write skew anomaly with concurrent transactions".
|
|
// which means we probably still only want one writeTx at a time
|
|
// even for badger.
|
|
writeCount int
|
|
writeByteCount int
|
|
|
|
o Txo
|
|
sn int64 // serial number
|
|
|
|
done bool
|
|
}
|
|
|
|
func (tx *BadgerTx) IsDone() (done bool) {
|
|
tx.mu.Lock()
|
|
done = tx.done
|
|
tx.mu.Unlock()
|
|
return
|
|
}
|
|
|
|
func (tx *BadgerTx) Group() *TxGroup {
|
|
return tx.o.Group
|
|
}
|
|
|
|
func (tx *BadgerTx) Options() Txo {
|
|
return tx.o
|
|
}
|
|
|
|
func (tx *BadgerTx) Type() string {
|
|
return BadgerTxn
|
|
}
|
|
|
|
func (tx *BadgerTx) UseRowCache() bool {
|
|
//the row cache speeds up queries.
|
|
return rbf.EnableRowCache
|
|
}
|
|
|
|
// overWriteOurAllocs provides detection of memory
|
|
// access outside the transactional context, similar to the
|
|
// old school electric fence techniques but without setting
|
|
// memory mappings to read-only... instead we just zero
|
|
// out the memory allocated to roaring containers by a
|
|
// transaction after the commit or rollback. This,
|
|
// hopefully, will cause some downstream confusion and
|
|
// test failures, which we can use to locate who has been
|
|
// holding on to memory they should have copied prior
|
|
// to transaction commit.
|
|
func (tx *BadgerTx) overWriteOurAllocs() {
|
|
|
|
tx.acMu.Lock()
|
|
defer tx.acMu.Unlock()
|
|
for _, s := range tx.ourAllocs {
|
|
|
|
// The Go compiler recognizes the following pattern and inserts
|
|
// an efficient memclr instruction.
|
|
// See https://github.com/golang/go/issues/5373
|
|
// and https://codereview.appspot.com/137880043
|
|
for i := range s {
|
|
s[i] = 0
|
|
// or
|
|
// Seebs suggested we might see even more crashes :)
|
|
// but since it will be slow (no memclr), we'll leave the default 0 for now.
|
|
//s[i] = -2
|
|
}
|
|
}
|
|
// keep this around if we need to activate out-of-mmap memory access again.
|
|
//for _, v := range tx.ourContainers {
|
|
//v.Invalid = true
|
|
//v.Tx = tx
|
|
//}
|
|
}
|
|
|
|
// Pointer gives us a memory address for the underlying transaction for debugging.
|
|
// It is public because we use it in roaring to report invalid container memory access
|
|
// outside of a transaction.
|
|
func (tx *BadgerTx) Pointer() string {
|
|
return fmt.Sprintf("%p", tx)
|
|
}
|
|
|
|
// Rollback rolls back the transaction.
|
|
func (tx *BadgerTx) Rollback() {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
//pp("BadgerTx.Rollback p=%p, its: '%v' initloc: '%v',\n rollbackloc:'%v'", tx, tx.Db.UnprotectedListOpenItAsString(), tx.initloc, stack())
|
|
tx.tx.Discard() // must hold tx.mu mutex lock
|
|
|
|
tx.Db.muOpenTxIt.Lock()
|
|
delete(tx.Db.openTx, tx)
|
|
tx.Db.muOpenTxIt.Unlock()
|
|
|
|
if tx.doAllocZero {
|
|
// and clear our allocs, to find code using them outside of a txn.
|
|
tx.overWriteOurAllocs()
|
|
}
|
|
|
|
tx.Db.unprotectedCleanupTx(tx)
|
|
}
|
|
|
|
// Commit commits the transaction to permanent storage.
|
|
// Commits can handle up to 100k updates to fragments
|
|
// at once, but not more. This is a BadgerDB imposed limit.
|
|
func (tx *BadgerTx) Commit() error {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
tx.Db.muOpenTxIt.Lock()
|
|
delete(tx.Db.openTx, tx)
|
|
tx.Db.muOpenTxIt.Unlock()
|
|
|
|
//pp("BadgerTx.Commit (write:%v) p=%p, stackID=%x openit: '%v' initloc: '%v', commitloc:\n%v", tx.write, tx, stackID, tx.Db.UnprotectedListOpenItAsString(), tx.initloc, stack())
|
|
|
|
err := tx.tx.CommitAt(1, nil) // must hold tx.mu mutex lock
|
|
|
|
if tx.doAllocZero {
|
|
tx.overWriteOurAllocs()
|
|
}
|
|
|
|
tx.Db.unprotectedCleanupTx(tx)
|
|
return err
|
|
}
|
|
|
|
// Readonly returns true iff the BadgerTx is read-only.
|
|
func (tx *BadgerTx) Readonly() bool {
|
|
return !tx.write
|
|
}
|
|
|
|
// RoaringBitmap returns the roaring.Bitmap for all bits in the fragment.
|
|
func (tx *BadgerTx) RoaringBitmap(index, field, view string, shard uint64) (*roaring.Bitmap, error) {
|
|
|
|
return tx.OffsetRange(index, field, view, shard, 0, 0, LeftShifted16MaxContainerKey)
|
|
}
|
|
|
|
// Container returns the requested roaring.Container, selected by fragment and ckey
|
|
func (tx *BadgerTx) Container(index, field, view string, shard uint64, ckey uint64) (c *roaring.Container, err error) {
|
|
|
|
// values returned from Get() are only valid while the transaction
|
|
// is open. If you need to use a value outside of the transaction then
|
|
// you must use copy() to copy it to another byte slice.
|
|
// BUT here we are already inside the Txn.
|
|
|
|
bkey := txkey.Key(index, field, view, shard, ckey)
|
|
tx.mu.Lock()
|
|
var item *badger.Item
|
|
item, err = tx.tx.Get(bkey)
|
|
tx.mu.Unlock()
|
|
if err == badger.ErrKeyNotFound {
|
|
// Seems crazy, but we, for now at least,
|
|
// match what RoaringTx does by returning nil, nil.
|
|
return nil, nil
|
|
} else {
|
|
panicOn(err)
|
|
}
|
|
|
|
err = item.Value(func(v []byte) error {
|
|
// This func with val would only be called if item.Value encounters no error
|
|
c = tx.toContainer(item.UserMeta(), v)
|
|
return nil
|
|
})
|
|
panicOn(err)
|
|
return
|
|
}
|
|
|
|
func (w *BadgerDBWrapper) NewWriteBatch() {
|
|
w.muDb.Lock()
|
|
defer w.muDb.Unlock()
|
|
if w.writeBatch != nil {
|
|
panic("must FlushWriteBatch() before calling NewWriteBatch()")
|
|
}
|
|
w.writeBatch = w.db.NewWriteBatch()
|
|
}
|
|
|
|
// Flush any remaining un-committed writes in progress.
|
|
func (w *BadgerDBWrapper) FlushWriteBatch() (err error) {
|
|
w.muDb.Lock()
|
|
defer w.muDb.Unlock()
|
|
if w.writeBatch == nil {
|
|
panic("CommitWriteBatch error: no batch in progress")
|
|
}
|
|
err = w.writeBatch.Flush()
|
|
w.writeBatch = nil
|
|
return
|
|
}
|
|
|
|
// Cancel any remaining un-committed writes in progress.
|
|
func (w *BadgerDBWrapper) CancelWriteBatch() {
|
|
w.muDb.Lock()
|
|
defer w.muDb.Unlock()
|
|
if w.writeBatch == nil {
|
|
panic("CancelWriteBatch error: no batch in progress")
|
|
}
|
|
w.writeBatch.Cancel()
|
|
}
|
|
|
|
// PutContainer stores rc under the specified fragment and container ckey.
|
|
func (tx *BadgerTx) PutContainer(index, field, view string, shard uint64, ckey uint64, rc *roaring.Container) error {
|
|
|
|
bkey := txkey.Key(index, field, view, shard, ckey)
|
|
var by []byte
|
|
|
|
ct := roaring.ContainerType(rc)
|
|
|
|
switch ct {
|
|
case roaring.ContainerArray:
|
|
by = fromArray16(roaring.AsArray(rc))
|
|
case roaring.ContainerBitmap:
|
|
by = fromArray64(roaring.AsBitmap(rc))
|
|
case roaring.ContainerRun:
|
|
by = fromInterval16(roaring.AsRuns(rc))
|
|
case roaring.ContainerNil:
|
|
panic("wat? nil roaring.Container is unexpected, no?!?")
|
|
default:
|
|
panic(fmt.Sprintf("unknown roaring.Container type: %v", ct))
|
|
}
|
|
|
|
entry := badger.NewEntry(bkey, by).WithMeta(ct)
|
|
|
|
tx.Db.muDb.Lock()
|
|
if tx.Db.writeBatch != nil {
|
|
err := tx.Db.writeBatch.SetEntry(entry) // Will create txns as needed.
|
|
tx.Db.muDb.Unlock()
|
|
return err
|
|
}
|
|
tx.Db.muDb.Unlock()
|
|
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
tx.writeCount++
|
|
sz := len(by) + len(bkey) + 2
|
|
tx.writeByteCount += sz
|
|
|
|
// The integration tests do large bit level loads that exceed 10MB.
|
|
// So we autocommit and start a new Txn if we are about to
|
|
// write too much into one Txn.
|
|
//
|
|
// The badger defaults limits are currently:
|
|
// maxBatchCount:104857, maxBatchSize:10066329
|
|
//
|
|
// However, emprirically we still get ErrTnTooBig when
|
|
// tx.writeByteCount=5884222; or when tx.writeCount=16197.
|
|
// So duck under both those thresholds by some margin.
|
|
if tx.writeCount > 100 || tx.writeByteCount > 2000000 {
|
|
// avoid ErrTxnTooBig by commiting before going over the limits,
|
|
// because then we get a error: "Transaction Conflict. Please retry."
|
|
err := tx.tx.CommitAt(1, nil)
|
|
panicOn(err)
|
|
// badger docs:
|
|
// `ErrConflict is returned when a transaction conflicts with another transaction. This can
|
|
// happen if the read rows had been updated concurrently by another transaction.
|
|
// ErrConflict = errors.New("Transaction Conflict. Please retry")`
|
|
//if err == badger.ErrConflict {
|
|
// problem is, we don't have the previous entry handy now.
|
|
//}
|
|
if err != nil {
|
|
panic(fmt.Sprintf("commit failed on bkey '%v': err '%v'", string(bkey), err))
|
|
}
|
|
tx.tx = tx.Db.db.NewTransactionAt(1, tx.write)
|
|
|
|
tx.writeCount = 1
|
|
tx.writeByteCount = sz
|
|
}
|
|
err := tx.tx.SetEntry(entry)
|
|
|
|
// ErrTxnTooBig is returned if too many writes are fit into a single transaction.
|
|
// badger docs: "An ErrTxnTooBig will be reported in case the number of pending
|
|
// writes/deletes in the transaction exceeds a certain limit. In that case, it
|
|
// is best to commit the transaction and start a new transaction immediately."
|
|
//
|
|
if err == badger.ErrTxnTooBig {
|
|
|
|
err = tx.tx.CommitAt(1, nil)
|
|
if err != nil {
|
|
panic(fmt.Sprintf("commit after TooBig failed on bkey '%v': err '%v'", string(bkey), err))
|
|
}
|
|
|
|
tx.tx = tx.Db.db.NewTransactionAt(1, tx.write)
|
|
//vv("NewTx write txn (p=%p) on gid=%v. b/c TooBig writeCount=%v; writeByteCount=%v", tx.tx, curGID(), tx.writeCount, tx.writeByteCount)
|
|
tx.writeCount = 1
|
|
tx.writeByteCount = sz
|
|
|
|
err = tx.tx.SetEntry(entry)
|
|
panicOn(err)
|
|
//panic(fmt.Sprintf("got error badger.ErrTxnTooBig, but we shoud never get this now; len(by) = %v; len(bkey)=%v; vs limit is 10MB. tx.writeCount=%v; tx.writeByteCount=%v;", len(by), len(bkey), tx.writeCount, tx.writeByteCount))
|
|
}
|
|
return err
|
|
}
|
|
|
|
// RemoveContainer deletes the container specified by the shard and container key ckey
|
|
func (tx *BadgerTx) RemoveContainer(index, field, view string, shard uint64, ckey uint64) error {
|
|
bkey := txkey.Key(index, field, view, shard, ckey)
|
|
tx.mu.Lock()
|
|
err := tx.tx.Delete(bkey)
|
|
tx.mu.Unlock()
|
|
return err
|
|
}
|
|
|
|
// Add sets all the a bits hot in the specified fragment.
|
|
func (tx *BadgerTx) Add(index, field, view string, shard uint64, batched bool, a ...uint64) (changeCount int, err error) {
|
|
return tx.addOrRemove(index, field, view, shard, batched, false, a...)
|
|
}
|
|
|
|
// Remove clears all the specified a bits in the chosen fragment.
|
|
func (tx *BadgerTx) Remove(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) {
|
|
const batched = false
|
|
const remove = true
|
|
return tx.addOrRemove(index, field, view, shard, batched, remove, a...)
|
|
}
|
|
|
|
func (tx *BadgerTx) addOrRemove(index, field, view string, shard uint64, batched, remove bool, a ...uint64) (changeCount int, err error) {
|
|
// pure hack to match RoaringTx
|
|
defer func() {
|
|
if !remove && !batched {
|
|
if changeCount > 0 {
|
|
changeCount = 1
|
|
}
|
|
}
|
|
}()
|
|
|
|
if len(a) == 0 {
|
|
return 0, nil
|
|
}
|
|
|
|
// have to sort, b/c input is not always sorted.
|
|
sort.Slice(a, func(i, j int) bool { return a[i] < a[j] })
|
|
|
|
var lastHi uint64 = math.MaxUint64 // highbits is always less than this starter.
|
|
var rc *roaring.Container
|
|
var hi uint64
|
|
var lo uint16
|
|
|
|
for i, v := range a {
|
|
|
|
hi, lo = highbits(v), lowbits(v)
|
|
if hi != lastHi {
|
|
// either first time through, or changed to a different container.
|
|
// do we need put the last updated container now?
|
|
if i > 0 {
|
|
// not first time through, write what we got.
|
|
if remove && (rc == nil || rc.N() == 0) {
|
|
err = tx.RemoveContainer(index, field, view, shard, lastHi)
|
|
panicOn(err)
|
|
} else {
|
|
err = tx.PutContainer(index, field, view, shard, lastHi, rc)
|
|
panicOn(err)
|
|
}
|
|
}
|
|
// get the next container
|
|
rc, err = tx.Container(index, field, view, shard, hi)
|
|
panicOn(err)
|
|
} // else same container, keep adding bits to rct.
|
|
chng := false
|
|
// rc can be nil before, and nil after, in both Remove/Add below.
|
|
// The roaring container add() and remove() methods handle this.
|
|
if remove {
|
|
rc, chng = rc.Remove(lo)
|
|
} else {
|
|
rc, chng = rc.Add(lo)
|
|
}
|
|
if chng {
|
|
changeCount++
|
|
}
|
|
lastHi = hi
|
|
}
|
|
// write the last updates.
|
|
if remove {
|
|
if rc == nil || rc.N() == 0 {
|
|
err = tx.RemoveContainer(index, field, view, shard, hi)
|
|
panicOn(err)
|
|
} else {
|
|
err = tx.PutContainer(index, field, view, shard, hi, rc)
|
|
panicOn(err)
|
|
}
|
|
} else {
|
|
if rc == nil || rc.N() == 0 {
|
|
panic("there should be no way to have an empty bitmap AFTER an Add() operation")
|
|
}
|
|
err = tx.PutContainer(index, field, view, shard, hi, rc)
|
|
panicOn(err)
|
|
}
|
|
return
|
|
}
|
|
|
|
// Contains returns exists true iff the bit chosen by key is
|
|
// hot (set to 1) in specified fragment.
|
|
func (tx *BadgerTx) Contains(index, field, view string, shard uint64, key uint64) (exists bool, err error) {
|
|
|
|
lo, hi := lowbits(key), highbits(key)
|
|
bkey := txkey.Key(index, field, view, shard, hi)
|
|
tx.mu.Lock()
|
|
item, err := tx.tx.Get(bkey)
|
|
tx.mu.Unlock()
|
|
if err == badger.ErrKeyNotFound {
|
|
return false, nil
|
|
}
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
err = item.Value(func(v []byte) error {
|
|
// This func with val would only be called if item.Value encounters no error
|
|
c := tx.toContainer(item.UserMeta(), v)
|
|
exists = c.Contains(lo)
|
|
return nil
|
|
})
|
|
return exists, err
|
|
}
|
|
|
|
func (tx *BadgerTx) SliceOfShards(index, field, view, optionalViewPath string) (sliceOfShards []uint64, err error) {
|
|
|
|
prefix := txkey.AllShardPrefix(index, field, view)
|
|
|
|
bi := NewBadgerIterator(tx, prefix)
|
|
defer bi.Close()
|
|
bi.Seek(prefix)
|
|
if !bi.it.Valid() {
|
|
return
|
|
}
|
|
lastShard := uint64(0)
|
|
firstDone := false
|
|
for bi.Next() {
|
|
item := bi.it.Item()
|
|
key := item.Key()
|
|
shard := txkey.ShardFromKey(key)
|
|
if firstDone {
|
|
if shard != lastShard {
|
|
sliceOfShards = append(sliceOfShards, shard)
|
|
}
|
|
lastShard = shard
|
|
} else {
|
|
// first time
|
|
lastShard = shard
|
|
firstDone = true
|
|
sliceOfShards = append(sliceOfShards, shard)
|
|
}
|
|
|
|
}
|
|
return
|
|
}
|
|
|
|
// key is the container key for the first roaring Container
|
|
// roaring docs: Iterator returns a ContainterIterator which *after* a call to Next(), a call to Value() will
|
|
// return the first container at or after key. found will be true if a
|
|
// container is found at key.
|
|
//
|
|
// BadgerTx notes: We auto-stop at the end of this shard, not going beyond.
|
|
func (tx *BadgerTx) ContainerIterator(index, field, view string, shard uint64, firstRoaringContainerKey uint64) (citer roaring.ContainerIterator, found bool, err error) {
|
|
|
|
// needle example: "idx:'i';fld:'f';vw:'v';shd:'00000000000000000000';key@00000000000000000000"
|
|
needle := txkey.Key(index, field, view, shard, firstRoaringContainerKey)
|
|
|
|
// prefix example: "idx:'i';fld:'f';vw:'v';shard:'00000000000000000000';key@"
|
|
prefix := txkey.Prefix(index, field, view, shard)
|
|
|
|
bi := NewBadgerIterator(tx, prefix)
|
|
bi.Seek(needle)
|
|
if !bi.it.Valid() {
|
|
return bi, false, nil
|
|
}
|
|
|
|
if !bi.it.ValidForPrefix(prefix) {
|
|
return bi, false, nil
|
|
}
|
|
item := bi.it.Item()
|
|
// have to compare b/c badger might give us valid iterator
|
|
// that is past our needle if needle isn't present.
|
|
return bi, bytes.Equal(item.Key(), needle), nil
|
|
}
|
|
|
|
// BadgerIterator is the iterator returned from a BadgerTx.ContainerIterator() call.
|
|
// It implements the roaring.ContainerIterator interface.
|
|
type BadgerIterator struct {
|
|
tx *BadgerTx
|
|
it *badger.Iterator
|
|
|
|
prefix []byte
|
|
seekto []byte
|
|
|
|
// seen counts how many Next() calls we have seen.
|
|
// It is used to match roaring.ContainerIterator semantics.
|
|
// Also useful for testing.
|
|
seen int
|
|
}
|
|
|
|
// NewBadgerIterator creates an iterator on tx that will
|
|
// only return txkey.Keys that start with prefix.
|
|
func NewBadgerIterator(tx *BadgerTx, prefix []byte) (bi *BadgerIterator) {
|
|
|
|
opts := badger.DefaultIteratorOptions
|
|
opts.PrefetchValues = false // else by default, pre-fetches the 1st 100 values, which would be slow.
|
|
opts.Reverse = false
|
|
|
|
tx.mu.Lock()
|
|
it := tx.tx.NewIterator(opts)
|
|
tx.mu.Unlock()
|
|
|
|
bi = &BadgerIterator{
|
|
tx: tx,
|
|
it: it,
|
|
prefix: prefix,
|
|
}
|
|
tx.Db.muOpenTxIt.Lock()
|
|
if tx.Db.openIt == nil {
|
|
tx.Db.openIt = make(map[*BadgerIterator]bool)
|
|
}
|
|
tx.Db.openIt[bi] = false // true for reverse, false for forward iteration.
|
|
tx.Db.muOpenTxIt.Unlock()
|
|
|
|
bi.it.Seek(prefix)
|
|
return
|
|
}
|
|
|
|
// NewBadgerReverseIterator makes a highest-to-lowest key iterator.
|
|
// Only keys that are prefixed with prefix will be returned.
|
|
// seekto tells where to start, and should be typically shard+1
|
|
// to start at the end of shard. Really only used in Max() at the moment.
|
|
// After creating a reverse badger iterator it, we will call it.Seek(seekto).
|
|
func NewBadgerReverseIterator(tx *BadgerTx, prefix, seekto []byte) (bi *BadgerIterator) {
|
|
|
|
tx.Db.muOpenTxIt.Lock()
|
|
defer tx.Db.muOpenTxIt.Unlock()
|
|
|
|
opts := badger.DefaultIteratorOptions
|
|
opts.PrefetchValues = false // else by default, pre-fetches the 1st 100 values, which would be slow.
|
|
opts.Reverse = true
|
|
opts.Prefix = prefix // possible storage IOPs optimization by badger
|
|
tx.mu.Lock()
|
|
it := tx.tx.NewIterator(opts)
|
|
tx.mu.Unlock()
|
|
|
|
bi = &BadgerIterator{
|
|
tx: tx,
|
|
it: it,
|
|
prefix: prefix,
|
|
seekto: seekto,
|
|
}
|
|
if tx.Db.openIt == nil {
|
|
tx.Db.openIt = make(map[*BadgerIterator]bool)
|
|
}
|
|
bi.tx.Db.openIt[bi] = true // true for reverse, false for forward iteration.
|
|
bi.it.Seek(seekto)
|
|
return
|
|
}
|
|
|
|
// Close tells the database and transaction that the user is done
|
|
// with the iterator.
|
|
// From the badger docs: It is important to call this when you're done with iteration.
|
|
// else you will get an error on tx.Discard()/Commit().
|
|
func (bi *BadgerIterator) Close() {
|
|
|
|
bi.tx.Db.muOpenTxIt.Lock()
|
|
delete(bi.tx.Db.openIt, bi)
|
|
bi.it.Close()
|
|
|
|
bi.tx.Db.muOpenTxIt.Unlock()
|
|
}
|
|
|
|
// Valid returns false if there are no more values in the iterator's range.
|
|
func (bi *BadgerIterator) Valid() bool {
|
|
return bi.it.Valid()
|
|
}
|
|
|
|
// Seek allows the iterator to start at needle instead of the global begining.
|
|
func (bi *BadgerIterator) Seek(needle []byte) {
|
|
bi.it.Seek(needle)
|
|
}
|
|
|
|
// Next advances the iterator.
|
|
func (bi *BadgerIterator) Next() bool {
|
|
|
|
// have to skip the first bi.it.Next() call because badger iterators point to the
|
|
// first value immediately, but Pilosa iterators must have Next() called
|
|
// on a fresh iterator to get the first value.
|
|
if bi.seen > 0 {
|
|
bi.it.Next()
|
|
}
|
|
bi.seen++
|
|
return bi.it.ValidForPrefix(bi.prefix) // does the bi.it.Valid() inside and false if not valid always.
|
|
}
|
|
|
|
// Value retrieves what is pointed at currently by the iterator.
|
|
func (bi *BadgerIterator) Value() (containerKey uint64, c *roaring.Container) {
|
|
if !bi.it.Valid() {
|
|
panic("bi.it not valid")
|
|
}
|
|
item := bi.it.Item()
|
|
if item == nil {
|
|
panic("item was nil")
|
|
}
|
|
key := item.Key()
|
|
containerKey = txkey.KeyExtractContainerKey(key)
|
|
|
|
err := item.Value(func(v []byte) error {
|
|
c = bi.tx.toContainer(item.UserMeta(), v)
|
|
return nil
|
|
})
|
|
panicOn(err)
|
|
return
|
|
}
|
|
|
|
// badgerFinder implements roaring.IteratorFinder.
|
|
// It is used by BadgerTx.ForEach()
|
|
type badgerFinder struct {
|
|
tx *BadgerTx
|
|
index string
|
|
field string
|
|
view string
|
|
shard uint64
|
|
needClose []Closer
|
|
}
|
|
|
|
// FindIterator lets badgerFinder implement the roaring.FindIterator interface.
|
|
func (bf *badgerFinder) FindIterator(seek uint64) (roaring.ContainerIterator, bool) {
|
|
a, found, err := bf.tx.ContainerIterator(bf.index, bf.field, bf.view, bf.shard, seek)
|
|
panicOn(err)
|
|
bf.needClose = append(bf.needClose, a)
|
|
return a, found
|
|
}
|
|
|
|
// Close closes all bf.needClose listed Closers.
|
|
func (bf *badgerFinder) Close() {
|
|
for _, i := range bf.needClose {
|
|
i.Close()
|
|
}
|
|
}
|
|
|
|
// NewTxIterator returns a *roaring.Iterator that MUST have Close() called on it BEFORE
|
|
// the transaction Commits or Rollsback.
|
|
func (tx *BadgerTx) NewTxIterator(index, field, view string, shard uint64) *roaring.Iterator {
|
|
bf := &badgerFinder{tx: tx, index: index, field: field, view: view, shard: shard, needClose: make([]Closer, 0)}
|
|
itr := roaring.NewIterator(bf)
|
|
return itr
|
|
}
|
|
|
|
// ForEach applies fn to each bitmap in the fragment.
|
|
func (tx *BadgerTx) ForEach(index, field, view string, shard uint64, fn func(i uint64) error) error {
|
|
itr := tx.NewTxIterator(index, field, view, shard)
|
|
defer itr.Close()
|
|
|
|
// Seek can create many container iterators, thus bf.Close() needClose list.
|
|
itr.Seek(0)
|
|
// v is the bit we are operating on.
|
|
for v, eof := itr.Next(); !eof; v, eof = itr.Next() {
|
|
if err := fn(v); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// ForEachRange applies fn on the selected range of bits on the chosen fragment.
|
|
func (tx *BadgerTx) ForEachRange(index, field, view string, shard uint64, start, end uint64, fn func(uint64) error) error {
|
|
|
|
itr := tx.NewTxIterator(index, field, view, shard)
|
|
defer itr.Close()
|
|
|
|
itr.Seek(start)
|
|
|
|
// v is the bit we are operating on.
|
|
for v, eof := itr.Next(); !eof && v < end; v, eof = itr.Next() {
|
|
if err := fn(v); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Count operates on the full bitmap level, so it sums over all the containers
|
|
// in the bitmap.
|
|
func (tx *BadgerTx) Count(index, field, view string, shard uint64) (uint64, error) {
|
|
|
|
a, found, err := tx.ContainerIterator(index, field, view, shard, 0)
|
|
panicOn(err)
|
|
defer a.Close()
|
|
if !found {
|
|
return 0, nil
|
|
}
|
|
result := int32(0)
|
|
for a.Next() {
|
|
ckey, cont := a.Value()
|
|
_ = ckey
|
|
result += cont.N()
|
|
}
|
|
|
|
return uint64(result), nil
|
|
}
|
|
|
|
// Max is the maximum bit-value in your bitmap.
|
|
// Returns zero if the bitmap is empty. Odd, but this is what roaring.Max does.
|
|
func (tx *BadgerTx) Max(index, field, view string, shard uint64) (uint64, error) {
|
|
|
|
prefix := txkey.Prefix(index, field, view, shard)
|
|
seekto := txkey.Prefix(index, field, view, shard+1)
|
|
|
|
it := NewBadgerReverseIterator(tx, prefix, seekto) // this iterator is still open, when we commit/discard tx.
|
|
defer it.Close()
|
|
|
|
if !it.it.Valid() {
|
|
return 0, nil
|
|
}
|
|
hb, rc := it.Value() // getting it returns invalid, as in empty iterator
|
|
lb := rc.Max()
|
|
|
|
return hb<<16 | uint64(lb), nil
|
|
}
|
|
|
|
// Min returns the smallest bit set in the fragment. If no bit is hot,
|
|
// the second return argument is false.
|
|
func (tx *BadgerTx) Min(index, field, view string, shard uint64) (uint64, bool, error) {
|
|
|
|
// Seek can create many container iterators, thus the bf.Close() needClose list.
|
|
bf := &badgerFinder{tx: tx, index: index, field: field, view: view, shard: shard, needClose: make([]Closer, 0)}
|
|
defer bf.Close()
|
|
itr := roaring.NewIterator(bf)
|
|
|
|
itr.Seek(0)
|
|
|
|
// v is the bit we are operating on.
|
|
v, eof := itr.Next()
|
|
if eof {
|
|
return 0, false, nil
|
|
}
|
|
return v, true, nil
|
|
}
|
|
|
|
// UnionInPlace unions all the others Bitmaps into a new Bitmap, and then writes it to the
|
|
// specified fragment.
|
|
func (tx *BadgerTx) UnionInPlace(index, field, view string, shard uint64, others ...*roaring.Bitmap) error {
|
|
|
|
rbm, err := tx.RoaringBitmap(index, field, view, shard)
|
|
panicOn(err)
|
|
|
|
rbm.UnionInPlace(others...)
|
|
// iterate over the containers that changed within rbm, and write them back to disk.
|
|
|
|
it, found := rbm.Containers.Iterator(0)
|
|
_ = found // don't care about the value of found, because first containerKey might be > 0
|
|
|
|
for it.Next() {
|
|
containerKey, rc := it.Value()
|
|
|
|
// TODO: only write the changed ones back, as optimization?
|
|
// Compare to ImportRoaringBits.
|
|
err := tx.PutContainer(index, field, view, shard, containerKey, rc)
|
|
panicOn(err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// CountRange returns the count of hot bits in the start, end range on the fragment.
|
|
// roaring.countRange counts the number of bits set between [start, end).
|
|
func (tx *BadgerTx) CountRange(index, field, view string, shard uint64, start, end uint64) (n uint64, err error) {
|
|
|
|
if tx.frag == nil {
|
|
return tx.countRangeNoFrag(index, field, view, shard, start, end)
|
|
}
|
|
|
|
// For speed, exploit the fact that on startup the rowCache will
|
|
// have already loaded fragments.
|
|
rowID := start / ShardWidth
|
|
row, err := tx.frag.unprotectedRow(tx, rowID)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
return row.Count(), nil
|
|
}
|
|
|
|
// CountRange returns the count of hot bits in the start, end range on the fragment.
|
|
// roaring.countRange counts the number of bits set between [start, end).
|
|
func (tx *BadgerTx) countRangeNoFrag(index, field, view string, shard uint64, start, end uint64) (n uint64, err error) {
|
|
|
|
if start >= end {
|
|
return 0, nil
|
|
}
|
|
|
|
skey := highbits(start)
|
|
ekey := highbits(end)
|
|
|
|
citer, found, err := tx.ContainerIterator(index, field, view, shard, skey)
|
|
_ = found
|
|
panicOn(err)
|
|
|
|
defer citer.Close()
|
|
|
|
// If range is entirely in one container then just count that range.
|
|
if skey == ekey {
|
|
citer.Next()
|
|
_, c := citer.Value()
|
|
return uint64(c.CountRange(int32(lowbits(start)), int32(lowbits(end)))), nil
|
|
}
|
|
|
|
for citer.Next() {
|
|
k, c := citer.Value()
|
|
if k < skey {
|
|
citer.Close()
|
|
panic(fmt.Sprintf("should be impossible for k(%v) to be less than skey(%v). tx p=%p", k, skey, tx))
|
|
}
|
|
|
|
// k > ekey handles the case when start > end and where start and end
|
|
// are in different containers. Same container case is already handled above.
|
|
if k > ekey {
|
|
break
|
|
}
|
|
if k == skey {
|
|
n += uint64(c.CountRange(int32(lowbits(start)), roaring.MaxContainerVal+1))
|
|
continue
|
|
}
|
|
if k < ekey {
|
|
n += uint64(c.N())
|
|
continue
|
|
}
|
|
if k == ekey {
|
|
n += uint64(c.CountRange(0, int32(lowbits(end))))
|
|
break
|
|
}
|
|
}
|
|
|
|
return n, nil
|
|
}
|
|
|
|
// OffsetRange creates a new roaring.Bitmap to return in other. For all the
|
|
// hot bits in [start, endx) of the chosen fragment, it stores
|
|
// them into other but with offset added to their bit position.
|
|
// The primary client is doing this, using ShardWidth, already; see
|
|
// fragment.rowFromStorage() in fragment.go. For example:
|
|
//
|
|
// data, err := tx.OffsetRange(f.index, f.field, f.view, f.shard,
|
|
// f.shard*ShardWidth, rowID*ShardWidth, (rowID+1)*ShardWidth)
|
|
// ^ offset ^ start ^ endx
|
|
//
|
|
// The start and endx arguments are container keys that have been shifted left by 16 bits;
|
|
// their highbits() will be taken to determine the actual container keys. This
|
|
// is done to conform to the roaring.OffsetRange() argument convention.
|
|
//
|
|
func (tx *BadgerTx) OffsetRange(index, field, view string, shard, offset, start, endx uint64) (other *roaring.Bitmap, err error) {
|
|
|
|
// roaring does these three checks in its OffsetRange
|
|
if lowbits(offset) != 0 {
|
|
panic("offset must not contain low bits")
|
|
}
|
|
if lowbits(start) != 0 {
|
|
panic("range start must not contain low bits")
|
|
}
|
|
if lowbits(endx) != 0 {
|
|
panic("range end must not contain low bits")
|
|
}
|
|
|
|
other = roaring.NewSliceBitmap()
|
|
off := highbits(offset)
|
|
hi0, hi1 := highbits(start), highbits(endx)
|
|
|
|
needle := txkey.Key(index, field, view, shard, hi0)
|
|
prefix := txkey.Prefix(index, field, view, shard)
|
|
|
|
n2, pre2 := txkey.KeyAndPrefix(index, field, view, shard, hi0)
|
|
if string(n2) != string(needle) {
|
|
panic(fmt.Sprintf("problem! n2(%v) != needle(%v), txkey.KeyAndPrefix not consitent with txkey.Key()", string(n2), string(needle)))
|
|
}
|
|
if string(pre2) != string(prefix) {
|
|
panic(fmt.Sprintf("problem! pre2(%v) != prefix(%v), txkey.KeyAndPrefix not consitent with txkey.Key()", string(pre2), string(prefix)))
|
|
}
|
|
|
|
it := NewBadgerIterator(tx, prefix)
|
|
defer it.Close()
|
|
it.Seek(needle)
|
|
for ; it.it.ValidForPrefix(prefix); it.Next() {
|
|
item := it.it.Item()
|
|
bkey := item.Key()
|
|
k := txkey.KeyExtractContainerKey(bkey)
|
|
|
|
// >= hi1 is correct b/c endx cannot have any lowbits set.
|
|
if uint64(k) >= hi1 {
|
|
break
|
|
}
|
|
destCkey := off + (k - hi0)
|
|
err := item.Value(func(v []byte) error {
|
|
c := tx.toContainer(item.UserMeta(), v)
|
|
other.Containers.Put(destCkey, c.Freeze())
|
|
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
return other, nil
|
|
}
|
|
|
|
// IncrementOpN increments the tx opcount by changedN
|
|
func (tx *BadgerTx) IncrementOpN(index, field, view string, shard uint64, changedN int) {
|
|
tx.opcount += changedN
|
|
}
|
|
|
|
// ImportRoaringBits handles deletes by setting clear=true.
|
|
// rowSet[rowID] returns the number of bit changed on that rowID.
|
|
func (tx *BadgerTx) ImportRoaringBits(index, field, view string, shard uint64, itr roaring.RoaringIterator, clear bool, log bool, rowSize uint64, data []byte) (changed int, rowSet map[uint64]int, err error) {
|
|
n := itr.Len()
|
|
if n == 0 {
|
|
return
|
|
}
|
|
rowSet = make(map[uint64]int)
|
|
|
|
var currRow uint64
|
|
|
|
var oldC *roaring.Container
|
|
for itrKey, synthC := itr.NextContainer(); synthC != nil; itrKey, synthC = itr.NextContainer() {
|
|
if rowSize != 0 {
|
|
currRow = itrKey / rowSize
|
|
}
|
|
nsynth := int(synthC.N())
|
|
if nsynth == 0 {
|
|
continue
|
|
}
|
|
// INVAR: nsynth > 0
|
|
|
|
oldC, err = tx.Container(index, field, view, shard, itrKey)
|
|
panicOn(err)
|
|
if err != nil {
|
|
return
|
|
}
|
|
|
|
if oldC == nil || oldC.N() == 0 {
|
|
// no container at the itrKey in badger (or all zero container).
|
|
if clear {
|
|
// changed of 0 and empty rowSet is perfect, no need to change the defaults.
|
|
continue
|
|
} else {
|
|
|
|
changed += nsynth
|
|
rowSet[currRow] += nsynth
|
|
|
|
err = tx.PutContainer(index, field, view, shard, itrKey, synthC)
|
|
if err != nil {
|
|
return
|
|
}
|
|
continue
|
|
}
|
|
}
|
|
|
|
if clear {
|
|
existN := oldC.N() // number of bits set in the old container
|
|
newC := oldC.Difference(synthC)
|
|
|
|
// update rowSet and changes
|
|
if newC.N() == existN {
|
|
// INVAR: do changed need adjusting? nope. same bit count,
|
|
// so no change could have happened.
|
|
continue
|
|
} else {
|
|
changes := int(existN - newC.N())
|
|
changed += changes
|
|
rowSet[currRow] -= changes
|
|
|
|
if tx.DeleteEmptyContainer && newC.N() == 0 {
|
|
err = tx.RemoveContainer(index, field, view, shard, itrKey)
|
|
if err != nil {
|
|
return
|
|
}
|
|
continue
|
|
}
|
|
err = tx.PutContainer(index, field, view, shard, itrKey, newC)
|
|
if err != nil {
|
|
return
|
|
}
|
|
continue
|
|
}
|
|
} else {
|
|
// setting bits
|
|
|
|
existN := oldC.N()
|
|
if existN == roaring.MaxContainerVal+1 {
|
|
// completely full container already, set will do nothing. so changed of 0 default is perfect.
|
|
continue
|
|
}
|
|
if existN == 0 {
|
|
// can nsynth be zero? No, because of the continue/invariant above where nsynth > 0
|
|
changed += nsynth
|
|
rowSet[currRow] += nsynth
|
|
err = tx.PutContainer(index, field, view, shard, itrKey, synthC)
|
|
if err != nil {
|
|
return
|
|
}
|
|
continue
|
|
}
|
|
|
|
newC := roaring.Union(oldC, synthC) // UnionInPlace was giving us crashes on overly large containers.
|
|
|
|
if roaring.ContainerType(newC) == roaring.ContainerBitmap {
|
|
newC.Repair() // update the bit-count so .n is valid. b/c UnionInPlace doesn't update it.
|
|
}
|
|
if newC.N() != existN {
|
|
changes := int(newC.N() - existN)
|
|
changed += changes
|
|
rowSet[currRow] += changes
|
|
|
|
err = tx.PutContainer(index, field, view, shard, itrKey, newC)
|
|
if err != nil {
|
|
panicOn(err)
|
|
return
|
|
}
|
|
continue
|
|
}
|
|
}
|
|
}
|
|
return
|
|
}
|
|
|
|
func (tx *BadgerTx) toContainer(typ byte, v []byte) (c *roaring.Container) {
|
|
|
|
if len(v) == 0 {
|
|
return nil
|
|
}
|
|
|
|
var w []byte
|
|
useRowCache := tx.UseRowCache()
|
|
if tx.doAllocZero || useRowCache {
|
|
// Do electric fence-inspired bad-memory read detection.
|
|
//
|
|
// The v []byte lives in BadgerDB's memory-mapped vlog-file,
|
|
// and Badger will recycle it after tx ends with rollback or commit.
|
|
//
|
|
// Problem is, at least some operations were not respecting transaction boundaries.
|
|
// This technique helped us find them. The rowCache was an example.
|
|
//
|
|
// See the global const DetectMemAccessPastTx
|
|
// at the top of txfactory.go to activate/deactivate this.
|
|
//
|
|
// Seebs suggested this nice variation: we could use individual mmaps for these
|
|
// copies, which would be unusable in production, but workable for testing, and then unmap them,
|
|
// which would get us probable segfaults on future accesses to them.
|
|
//
|
|
// The go runtime also has an -efence flag which may be similarly useful if really pressed.
|
|
//
|
|
w = make([]byte, len(v))
|
|
copy(w, v)
|
|
|
|
if !useRowCache {
|
|
// register w so we can catch out-of-tx memory access
|
|
tx.acMu.Lock()
|
|
defer tx.acMu.Unlock()
|
|
tx.ourAllocs = append(tx.ourAllocs, w)
|
|
}
|
|
} else {
|
|
w = v
|
|
}
|
|
|
|
switch typ {
|
|
case roaring.ContainerArray:
|
|
c = roaring.NewContainerArray(toArray16(w))
|
|
if tx.doAllocZero {
|
|
// tx.acMu was acquired above, and Unlock deferred.
|
|
tx.ourContainers = append(tx.ourContainers, c)
|
|
}
|
|
case roaring.ContainerBitmap:
|
|
c = roaring.NewContainerBitmap(-1, toArray64(w))
|
|
if tx.doAllocZero {
|
|
// tx.acMu was acquired above, and Unlock deferred.
|
|
tx.ourContainers = append(tx.ourContainers, c)
|
|
}
|
|
case roaring.ContainerRun:
|
|
c = roaring.NewContainerRun(toInterval16(w))
|
|
if tx.doAllocZero {
|
|
// tx.acMu was acquired above, and Unlock deferred.
|
|
tx.ourContainers = append(tx.ourContainers, c)
|
|
}
|
|
default:
|
|
panic(fmt.Sprintf("unknown container: %v", typ))
|
|
}
|
|
c.SetMapped(true)
|
|
return c
|
|
}
|
|
|
|
// StringifiedBadgerKeys returns a string with all the container
|
|
// keys available in badger.
|
|
func (w *BadgerDBWrapper) StringifiedBadgerKeys(optionalUseThisTx Tx, short bool) (r string) {
|
|
if optionalUseThisTx == nil {
|
|
tx, _ := w.NewTx(!writable, "<StringifiedBadgerKeys>", Txo{})
|
|
defer tx.Rollback()
|
|
r = stringifiedBadgerKeysTx(tx.(*BadgerTx), short)
|
|
return
|
|
}
|
|
|
|
btx, ok := optionalUseThisTx.(*BadgerTx)
|
|
if !ok {
|
|
return fmt.Sprintf("<not-a-BadgerTx-in-StringifiedBadgerKeys-was-%T>", optionalUseThisTx)
|
|
}
|
|
r = stringifiedBadgerKeysTx(btx, short)
|
|
return
|
|
}
|
|
|
|
// countBitsSet returns the number of bits set (or "hot") in
|
|
// the roaring container value found by the txkey.Key()
|
|
// formatted bkey.
|
|
func (tx *BadgerTx) countBitsSet(bkey []byte) (n int) {
|
|
|
|
item, err := tx.tx.Get(bkey)
|
|
if err == badger.ErrKeyNotFound {
|
|
panic(fmt.Sprintf("badger did not have value for bkey = '%v'", string(bkey)))
|
|
}
|
|
panicOn(err)
|
|
|
|
var rc *roaring.Container
|
|
err = item.Value(func(v []byte) error {
|
|
// This func with val would only be called if item.Value encounters no error
|
|
rc = tx.toContainer(item.UserMeta(), v)
|
|
return nil
|
|
})
|
|
panicOn(err)
|
|
|
|
n = int(rc.N())
|
|
return
|
|
}
|
|
|
|
func (tx *BadgerTx) Dump(short bool, shard uint64) {
|
|
fmt.Printf("BadgerTx %p Dump: %v\n", tx, stringifiedBadgerKeysTx(tx, short))
|
|
}
|
|
|
|
// stringifiedBadgerKeysTx reports all the badger keys and a
|
|
// corresponding blake3 hash viewable by txn within the entire
|
|
// badger database.
|
|
// It also reports how many bits are hot in the roaring container
|
|
// (how many bits are set, or 1 rather than 0).
|
|
//
|
|
// By convention, we must return the empty string if there
|
|
// are no keys present. The tests use this to confirm
|
|
// an empty database.
|
|
func stringifiedBadgerKeysTx(tx *BadgerTx, short bool) (r string) {
|
|
|
|
r = "allkeys:[\n"
|
|
it := tx.tx.NewIterator(badger.DefaultIteratorOptions) // PrefetchValues true okay here.
|
|
defer it.Close()
|
|
any := false
|
|
for it.Rewind(); it.Valid(); it.Next() {
|
|
any = true
|
|
item := it.Item()
|
|
bkey := item.Key()
|
|
key := txkey.ToString(bkey)
|
|
ckey := txkey.KeyExtractContainerKey(bkey)
|
|
h := ""
|
|
srbm := ""
|
|
err := item.Value(func(val []byte) error {
|
|
h = hash.Blake3sum16(val)
|
|
ct := tx.toContainer(item.UserMeta(), val)
|
|
cts := roaring.NewSliceContainers()
|
|
cts.Put(ckey, ct)
|
|
rbm := &roaring.Bitmap{Containers: cts}
|
|
srbm = BitmapAsString(rbm)
|
|
return nil
|
|
})
|
|
panicOn(err)
|
|
r += fmt.Sprintf("%v -> %v (%v hot)\n", key, h, tx.countBitsSet(bkey))
|
|
if !short {
|
|
r += " ......." + srbm + "\n"
|
|
}
|
|
}
|
|
r += "]\n all-in-blake3:" + hash.Blake3sum16([]byte(r))
|
|
|
|
if !any {
|
|
return ""
|
|
}
|
|
return "badger-" + r
|
|
}
|
|
|
|
func (w *BadgerDBWrapper) DeleteField(index, field, fieldPath string) error {
|
|
|
|
// under blue-green roaring_badger, the directory will not be found, b/c roaring will have
|
|
// already done the os.RemoveAll(). BUT, RemoveAll returns nil error in this case. Docs:
|
|
// "If the path does not exist, RemoveAll returns nil (no error)"
|
|
err := os.RemoveAll(fieldPath)
|
|
if err != nil {
|
|
return errors.Wrap(err, "removing directory")
|
|
}
|
|
prefix := txkey.FieldPrefix(index, field)
|
|
return w.DeletePrefix(prefix)
|
|
}
|
|
|
|
func (w *BadgerDBWrapper) DeleteFragment(index, field, view string, shard uint64, frag interface{}) error {
|
|
prefix := txkey.Prefix(index, field, view, shard)
|
|
return w.DeletePrefix(prefix)
|
|
}
|
|
|
|
func (w *BadgerDBWrapper) DeletePrefix(prefix []byte) error {
|
|
w.muDb.Lock()
|
|
defer w.muDb.Unlock()
|
|
|
|
// a) do key-ony iteration, no value fetch;
|
|
//
|
|
// b) do deletes in large batches, to avoid alot of txn overhead;
|
|
// per recommendation https://github.com/dgraph-io/badger/issues/598
|
|
//
|
|
// c) we do not, at present, try to maintain one large
|
|
// transaction with all the keys in a index in it. Because
|
|
// there can be too many keys. Hence the index will disappear
|
|
// in chucks of 100K keys, not atomically-all-at-once.
|
|
|
|
noMoreKeysWithPrefix := false
|
|
const maxDeletesPerTxn = 100000
|
|
|
|
for !noMoreKeysWithPrefix {
|
|
|
|
txn := w.db.NewTransactionAt(1, true)
|
|
//defer txn.Discard()
|
|
|
|
o := badger.DefaultIteratorOptions
|
|
o.AllVersions = false
|
|
o.PrefetchValues = false // key-only iteration, no values.
|
|
|
|
// note: panic: Unclosed iterator at time of Txn.Discard ? panic on segfault here?
|
|
// This means we messed up and Closed() the Database already; too early. For
|
|
// example in TxFactor.CloseIndex() in txfactory.go:331.
|
|
it := txn.NewIterator(o)
|
|
|
|
//defer it.Close()
|
|
n := 0
|
|
goners := make([][]byte, 0, maxDeletesPerTxn)
|
|
for it.Seek(prefix); it.ValidForPrefix(prefix); it.Next() {
|
|
|
|
// KeyCopy() is required; Key() means corruption and possible segfault.
|
|
key := it.Item().KeyCopy(nil)
|
|
goners = append(goners, key)
|
|
n++
|
|
if n >= maxDeletesPerTxn {
|
|
break
|
|
}
|
|
}
|
|
if !it.ValidForPrefix(prefix) {
|
|
noMoreKeysWithPrefix = true // done with the full delete of up to maxDeletesPerTxn
|
|
}
|
|
for _, key := range goners {
|
|
if err := txn.Delete(key); err != nil {
|
|
it.Close()
|
|
txn.Discard()
|
|
return err
|
|
}
|
|
}
|
|
it.Close()
|
|
err := txn.CommitAt(1, nil)
|
|
|
|
// err back from Update can be ErrConflict in case of
|
|
// a conflict. Badger docs: "Depending on the state
|
|
// of your application, you have the option to
|
|
// retry the operation if you receive this error."
|
|
panicOn(err)
|
|
|
|
} // end for: proceed to next bath of 100K keys
|
|
|
|
// Finally, run a garbage collection to delete values from the value log.
|
|
//
|
|
// "Only one GC is allowed at a time. If another value log GC
|
|
// is running, or DB has been closed, this would return an ErrRejected."
|
|
// -- https://godoc.org/github.com/dgraph-io/badger#DB.RunValueLogGC
|
|
// Still, we don't see a mutex inside the RunValueLogGC code, so
|
|
// lock muGC just to be sure.
|
|
w.muGC.Lock()
|
|
defer w.muGC.Unlock()
|
|
_ = w.db.RunValueLogGC(0.5)
|
|
|
|
return nil
|
|
}
|
|
|
|
func (tx *BadgerTx) RoaringBitmapReader(index, field, view string, shard uint64, fragmentPathForRoaring string) (r io.ReadCloser, sz int64, err error) {
|
|
|
|
rbm, err := tx.RoaringBitmap(index, field, view, shard)
|
|
if err != nil {
|
|
return nil, -1, errors.Wrap(err, "RoaringBitmapReader RoaringBitmap")
|
|
}
|
|
var buf bytes.Buffer
|
|
sz, err = rbm.WriteTo(&buf)
|
|
if err != nil {
|
|
return nil, -1, errors.Wrap(err, "RoaringBitmapReader rbm.WriteTo(buf)")
|
|
}
|
|
return ioutil.NopCloser(&buf), sz, err
|
|
}
|
|
|
|
func (w *BadgerDBWrapper) OpenListString() (r string) {
|
|
return "BadgerDBWrapper.OpenListString not yet implemented"
|
|
}
|
|
func (w *BadgerDBWrapper) OpenSnList() (sns []int64) {
|
|
return nil // panic("BadgerDBWrapper.OpenSnList not yet implemented")
|
|
}
|
|
|
|
// Sn retreives the serial number of the Tx.
|
|
func (tx *BadgerTx) Sn() int64 {
|
|
return tx.sn
|
|
}
|
|
|
|
func (w *BadgerDBWrapper) Update(fn func(txn *badger.Txn) error) error {
|
|
txn := w.db.NewTransactionAt(1, true)
|
|
defer txn.Discard()
|
|
|
|
if err := fn(txn); err != nil {
|
|
return err
|
|
}
|
|
|
|
return txn.CommitAt(1, nil)
|
|
}
|