mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-07 11:27:50 +00:00
Merge branch 'master' into pilosa-id-gen
This commit is contained in:
commit
a4f538ffec
29 changed files with 708 additions and 3771 deletions
|
|
@ -208,7 +208,7 @@ workflows:
|
|||
- setup
|
||||
matrix:
|
||||
parameters:
|
||||
test_make_target: ["test-race", "test-txstore-rbf", "test-txstore-rbf_lmdb"]
|
||||
test_make_target: ["test-race", "test-txstore-rbf", "test-txstore-rbf_bolt"]
|
||||
- test:
|
||||
name: test-shardwidth-22
|
||||
shard_width: "22"
|
||||
|
|
|
|||
6
Makefile
6
Makefile
|
|
@ -1,4 +1,4 @@
|
|||
.PHONY: build check-clean clean build-lattice cover cover-viz default docker docker-build docker-test docker-tag-push generate generate-protoc generate-pql generate-statik gometalinter install install-build-deps install-golangci-lint install-gometalinter install-protoc install-protoc-gen-gofast install-peg install-statik prerelease prerelease-upload release release-build test testv testv-race testvsub testvsub-race test-txstore-rbf_lmdb test-txstore-rbf
|
||||
.PHONY: build check-clean clean build-lattice cover cover-viz default docker docker-build docker-test docker-tag-push generate generate-protoc generate-pql generate-statik gometalinter install install-build-deps install-golangci-lint install-gometalinter install-protoc install-protoc-gen-gofast install-peg install-statik prerelease prerelease-upload release release-build test testv testv-race testvsub testvsub-race test-txstore-rbf
|
||||
|
||||
CLONE_URL=github.com/pilosa/pilosa
|
||||
MOD_VERSION=v2
|
||||
|
|
@ -309,6 +309,6 @@ install-gometalinter:
|
|||
test-txstore-rbf:
|
||||
PILOSA_TXSRC=rbf $(MAKE) testv-race
|
||||
|
||||
test-txstore-rbf_lmdb:
|
||||
PILOSA_TXSRC=rbf_lmdb $(MAKE) testv-race
|
||||
test-txstore-rbf_bolt:
|
||||
PILOSA_TXSRC=rbf_bolt $(MAKE) testv-race
|
||||
|
||||
|
|
|
|||
17
bolt.go
17
bolt.go
|
|
@ -37,6 +37,8 @@ import (
|
|||
bolt "go.etcd.io/bbolt"
|
||||
)
|
||||
|
||||
const isDebugRun = false
|
||||
|
||||
// boltRegistrar facilitates shutdown
|
||||
// of all the bolt databases started under
|
||||
// tests. Its needed because most tests don't cleanup
|
||||
|
|
@ -1419,6 +1421,21 @@ func (tx *BoltTx) toContainer(typ byte, v []byte) (r *roaring.Container) {
|
|||
return ToContainer(typ, w)
|
||||
}
|
||||
|
||||
func ToContainer(typ byte, w []byte) (c *roaring.Container) {
|
||||
switch typ {
|
||||
case roaring.ContainerArray:
|
||||
c = roaring.NewContainerArray(toArray16(w))
|
||||
case roaring.ContainerBitmap:
|
||||
c = roaring.NewContainerBitmap(-1, toArray64(w))
|
||||
case roaring.ContainerRun:
|
||||
c = roaring.NewContainerRun(toInterval16(w))
|
||||
default:
|
||||
panic(fmt.Sprintf("unknown container: %v", typ))
|
||||
}
|
||||
c.SetMapped(true)
|
||||
return c
|
||||
}
|
||||
|
||||
// StringifiedBoltKeys returns a string with all the container
|
||||
// keys available in bolt.
|
||||
func (w *BoltWrapper) StringifiedBoltKeys(optionalUseThisTx Tx, short bool) (r string) {
|
||||
|
|
|
|||
19
const_amd64.go
Normal file
19
const_amd64.go
Normal file
|
|
@ -0,0 +1,19 @@
|
|||
// 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.
|
||||
|
||||
// +build amd64
|
||||
|
||||
package pilosa
|
||||
|
||||
const TxInitialMmapSize = 4 << 30 // 4GB
|
||||
21
const_other.go
Normal file
21
const_other.go
Normal file
|
|
@ -0,0 +1,21 @@
|
|||
// 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.
|
||||
|
||||
// +build !amd64
|
||||
|
||||
package pilosa
|
||||
|
||||
// this is a stubbed out file to let 386/arm build.
|
||||
|
||||
const TxInitialMmapSize = 1 << 30 // 1GB
|
||||
|
|
@ -478,7 +478,6 @@ func (dbs *DBShard) DumpAll() {
|
|||
switch ty {
|
||||
case roaringTxn:
|
||||
case rbfTxn:
|
||||
case lmdbTxn:
|
||||
case boltTxn:
|
||||
default:
|
||||
panic(fmt.Sprintf("unknown txtyp: '%v'", ty))
|
||||
|
|
@ -605,8 +604,6 @@ func (per *DBPerShard) GetDBShard(index string, shard uint64, idx *Index) (dbs *
|
|||
registry = globalRoaringReg
|
||||
case rbfTxn:
|
||||
registry = globalRbfDBReg
|
||||
case lmdbTxn:
|
||||
registry = globalLMDBReg
|
||||
case boltTxn:
|
||||
registry = globalBoltReg
|
||||
default:
|
||||
|
|
@ -888,9 +885,9 @@ func listDirUnderDir(root string, includeRoot bool, requiredSuffix string, ignor
|
|||
// The blue is the destination -- this is always types[0].
|
||||
// The green source is always types[1]. The mnemonic is blue_geen.
|
||||
// The blue is first, so it is in types[0]. The green
|
||||
// is second, in types[1]. For example, with PILOSA_TXSRC=lmdb_roaring
|
||||
// we have lmdb as blue, and roaring as green. The contents of
|
||||
// lmdb must be empty or exactly match roaring. If lmdb
|
||||
// is second, in types[1]. For example, with PILOSA_TXSRC=bolt_roaring
|
||||
// we have bolt as blue, and roaring as green. The contents of
|
||||
// bolt must be empty or exactly match roaring. If bolt
|
||||
// starts empty, it will be populated from roaring by
|
||||
// populateBlueFromGreen().
|
||||
//
|
||||
|
|
|
|||
|
|
@ -74,7 +74,7 @@ func Test_DBPerShard_GetShardsForIndex_LocalOnly(t *testing.T) {
|
|||
orig := os.Getenv("PILOSA_TXSRC")
|
||||
defer os.Setenv("PILOSA_TXSRC", orig) // must restore or will mess up other tests!
|
||||
|
||||
for _, src := range []string{"lmdb", "roaring", "bolt", "rbf"} {
|
||||
for _, src := range []string{"roaring", "bolt", "rbf"} {
|
||||
|
||||
os.Setenv("PILOSA_TXSRC", src)
|
||||
|
||||
|
|
@ -131,14 +131,6 @@ rick/_exists/views/standard/fragments/217
|
|||
rick/_exists/views/standard/fragments/93
|
||||
rick/_exists/views/standard/fragments/219
|
||||
rick/_exists/views/standard/fragments/223
|
||||
`,
|
||||
"lmdb": `
|
||||
rick.index.txstores@@@/store-lmdb@@/shard.0093-lmdb@
|
||||
rick.index.txstores@@@/store-lmdb@@/shard.0215-lmdb@
|
||||
rick.index.txstores@@@/store-lmdb@@/shard.0217-lmdb@
|
||||
rick.index.txstores@@@/store-lmdb@@/shard.0219-lmdb@
|
||||
rick.index.txstores@@@/store-lmdb@@/shard.0221-lmdb@
|
||||
rick.index.txstores@@@/store-lmdb@@/shard.0223-lmdb@
|
||||
`,
|
||||
"bolt": `
|
||||
rick.index.txstores@@@/store-boltdb@@/shard.0093-boltdb@/bolt.db
|
||||
|
|
@ -173,13 +165,6 @@ func makeSampleRoaringDir(root, txsrc string, minBytes int, h *Holder) {
|
|||
shard = shards[i]
|
||||
}
|
||||
switch txsrc {
|
||||
case "lmdb":
|
||||
makeLMDBtestDB(root+sep+fn, h, shard)
|
||||
// also have to make the DBShard in our in-memory tree,
|
||||
// or else the search won't find it because
|
||||
// DBPerShard won't know anything about it.
|
||||
helperCreateDBShard(h, index, shard)
|
||||
continue
|
||||
case "bolt":
|
||||
makeBolttestDB(root+sep+fn, h, shard)
|
||||
helperCreateDBShard(h, index, shard)
|
||||
|
|
@ -210,14 +195,6 @@ func helperCreateDBShard(h *Holder, index string, shard uint64) {
|
|||
_ = dbs
|
||||
}
|
||||
|
||||
func makeLMDBtestDB(path string, h *Holder, shard uint64) {
|
||||
i := uint64(1)
|
||||
w, _ := mustOpenEmptyLMDBWrapper(path)
|
||||
LMDBMustSetBitvalue(w, "index", "field", "view", shard, i)
|
||||
w.Close()
|
||||
|
||||
}
|
||||
|
||||
func makeBolttestDB(path string, h *Holder, shard uint64) {
|
||||
i := uint64(1)
|
||||
w, _ := mustOpenEmptyBoltWrapper(path)
|
||||
|
|
|
|||
|
|
@ -17,7 +17,6 @@ package pilosa_test
|
|||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"reflect"
|
||||
"testing"
|
||||
|
||||
|
|
@ -28,59 +27,6 @@ import (
|
|||
"github.com/pilosa/pilosa/v2/test"
|
||||
)
|
||||
|
||||
func skipForNonLMDB(t *testing.T) {
|
||||
src := os.Getenv("PILOSA_TXSRC")
|
||||
if src != "lmdb" {
|
||||
t.Skip("skip if not lmdb")
|
||||
}
|
||||
}
|
||||
|
||||
var _ = skipForNonLMDB // happy linter
|
||||
|
||||
// Can't write it all to one shard like we do (did).
|
||||
func Test_DBPerShard_multiple_shards_used(t *testing.T) {
|
||||
skipForNonLMDB(t)
|
||||
c := test.MustRunCluster(t, 1)
|
||||
defer c.Close()
|
||||
hldr := c.GetHolder(0)
|
||||
index := "i"
|
||||
hldr.SetBit(index, "general", 10, 0)
|
||||
hldr.SetBit(index, "general", 10, ShardWidth+1)
|
||||
hldr.SetBit(index, "general", 10, ShardWidth+2)
|
||||
|
||||
hldr.SetBit(index, "general", 11, 2)
|
||||
hldr.SetBit(index, "general", 11, ShardWidth+2)
|
||||
|
||||
types := pilosa.MustTxsrcToTxtype("lmdb")
|
||||
idx := hldr.Index(index)
|
||||
shardsU := []uint64{0, 1, 2}
|
||||
pathShard := []string{}
|
||||
|
||||
// check that 3 different shard databases/files were made
|
||||
for i := 0; i < 2; i++ {
|
||||
|
||||
path, err := hldr.Txf().GetDBShardPath(index, shardsU[i], idx, types[0], !writable)
|
||||
panicOn(err)
|
||||
pathShard = append(pathShard, path)
|
||||
|
||||
if !DirExists(pathShard[i]) {
|
||||
panic(fmt.Sprintf("no shard made for pathShard[%v]='%v'", i, pathShard[i]))
|
||||
}
|
||||
sz, err := pilosa.DiskUse(pathShard[i], "")
|
||||
panicOn(err)
|
||||
|
||||
if sz < 100 {
|
||||
panic(fmt.Sprintf("shard %v was too small", i))
|
||||
}
|
||||
}
|
||||
|
||||
if res, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{Index: index, Query: `Union(Row(general=10), Row(general=11))`}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if columns := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, []uint64{0, 2, ShardWidth + 1, ShardWidth + 2}) {
|
||||
t.Fatalf("unexpected columns: %+v", columns)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAPI_SimplerOneNode_ImportColumnKey(t *testing.T) {
|
||||
|
||||
c := test.MustRunCluster(t, 1,
|
||||
|
|
|
|||
|
|
@ -5110,9 +5110,6 @@ func TestImportValueConcurrent(t *testing.T) {
|
|||
"blueGreenTx because the lack of transactional consistency " +
|
||||
"from Roaring-per-file will create false comparison " +
|
||||
"failures."))
|
||||
case lmdbTxn:
|
||||
t.Skip(fmt.Sprintf("skipping TestImportValueConcurrent under " +
|
||||
"lmdb since only a single writer is allowed at once."))
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -357,6 +357,10 @@ func (h *Handler) collectStats(next http.Handler) http.Handler {
|
|||
})
|
||||
}
|
||||
|
||||
// latticeRoutes lists the frontend routes that do not directly correspond to
|
||||
// backend routes, and require special handling.
|
||||
var latticeRoutes = []string{"/tables", "/query"}
|
||||
|
||||
// newRouter creates a new mux http router.
|
||||
func newRouter(handler *Handler) http.Handler {
|
||||
router := mux.NewRouter()
|
||||
|
|
@ -446,10 +450,12 @@ func newRouter(handler *Handler) http.Handler {
|
|||
latticeHandler := NewStatikHandler(handler)
|
||||
router.PathPrefix("/static").Handler(latticeHandler)
|
||||
router.Path("/").Handler(latticeHandler)
|
||||
router.Path("/vds").Handler(latticeHandler)
|
||||
router.Path("/favicon.png").Handler(latticeHandler)
|
||||
router.Path("/favicon.svg").Handler(latticeHandler)
|
||||
router.Path("/manifest.json").Handler(latticeHandler)
|
||||
for _, route := range latticeRoutes {
|
||||
router.Path(route).Handler(latticeHandler)
|
||||
}
|
||||
|
||||
router.Use(handler.queryArgValidator)
|
||||
router.Use(handler.addQueryContext)
|
||||
|
|
@ -522,11 +528,13 @@ func (s statikHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||
return
|
||||
}
|
||||
|
||||
// /vds is a front-end route, not a backend route. Without this check, refreshing at /vds
|
||||
// Without this check, refreshing the UI at e.g. /query
|
||||
// will request a nonexistent resource and return 404.
|
||||
if r.URL.String() == "/vds" {
|
||||
url, _ := url.Parse("/")
|
||||
r.URL = url
|
||||
for _, route := range latticeRoutes {
|
||||
if r.URL.String() == route {
|
||||
url, _ := url.Parse("/")
|
||||
r.URL = url
|
||||
}
|
||||
}
|
||||
|
||||
http.FileServer(s.statikFS).ServeHTTP(w, r)
|
||||
|
|
|
|||
456
lmdb_other.go
456
lmdb_other.go
|
|
@ -1,456 +0,0 @@
|
|||
// 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.
|
||||
|
||||
// +build !amd64
|
||||
|
||||
package pilosa
|
||||
|
||||
// this is a stubbed out file to let 386 build. lmdb won't work well
|
||||
// on 32-bit; not enough memory map address space.
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg"
|
||||
"github.com/pilosa/pilosa/v2/roaring"
|
||||
)
|
||||
|
||||
var _ = time.Now
|
||||
|
||||
const isDebugRun = false
|
||||
const TxInitialMmapSize = 1 << 30 // 1GB
|
||||
|
||||
func ToContainer(typ byte, w []byte) (c *roaring.Container) {
|
||||
panic("ToContainer not implemented yet on non-amd64")
|
||||
}
|
||||
|
||||
// lmdbRegistrar facilitates shutdown
|
||||
// of all the lmdb databases started under
|
||||
// tests. Its needed because most tests don't cleanup
|
||||
// the *Index(es) they create. But we still
|
||||
// want to shutdown lmdbDB 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 lmdbRegistrar struct {
|
||||
mu sync.Mutex
|
||||
mp map[*LMDBWrapper]bool
|
||||
|
||||
path2db map[string]*LMDBWrapper
|
||||
}
|
||||
|
||||
var globalLMDBReg *lmdbRegistrar = newLMDBTestRegistrar()
|
||||
|
||||
func newLMDBTestRegistrar() *lmdbRegistrar {
|
||||
|
||||
return &lmdbRegistrar{
|
||||
mp: make(map[*LMDBWrapper]bool),
|
||||
path2db: make(map[string]*LMDBWrapper),
|
||||
}
|
||||
}
|
||||
|
||||
func (r *lmdbRegistrar) OpenDBWrapper(path string, doAllocZero bool, rbfcfg *rbfcfg.Config) (DBWrapper, error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (r *lmdbRegistrar) Size() int {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// register each lmdb created under tests, so we
|
||||
// can clean them up. This is called by openLMDBWrapper() 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 *lmdbRegistrar) unprotectedRegister(w *LMDBWrapper) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// unregister removes w from r
|
||||
func (r *lmdbRegistrar) unregister(w *LMDBWrapper) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func DumpAllLMDB() {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// lmdbPath is a helper for determining the full directory
|
||||
// in which the lmdb database will be stored.
|
||||
func lmdbPath(path string) string {
|
||||
if !strings.HasSuffix(path, "-lmdb") {
|
||||
return path + "-lmdb"
|
||||
}
|
||||
return path
|
||||
}
|
||||
|
||||
// openLMDBDB opens the database in the bpath directoy
|
||||
// without deleting any prior content. Any LMDBDB
|
||||
// database directory will have the "-lmdb" suffix.
|
||||
//
|
||||
// openLMDBDB 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 lmdbDB
|
||||
// per bpath in this pilosa node.
|
||||
func (r *lmdbRegistrar) openLMDBWrapper(path0 string) (*LMDBWrapper, error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
var ErrShutdown = fmt.Errorf("shutting down")
|
||||
|
||||
// DeleteIndex deletes all the containers associated with
|
||||
// the named index from the lmdb database.
|
||||
func (w *LMDBWrapper) DeleteIndex(indexName string) error {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// statically confirm that LMDBTx satisfies the Tx interface.
|
||||
var _ Tx = (*LMDBTx)(nil)
|
||||
|
||||
// LMDBWrapper provides the NewLMDBTx() method.
|
||||
// Execute lmdbJob's via LMDBWrapper.submit(); these must
|
||||
// be done by the lmdb goroutine worker pool.
|
||||
type LMDBWrapper struct{}
|
||||
|
||||
func (w *LMDBWrapper) IsClosed() bool {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// NewLMDBTx produces LMDB 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 LMDBTx 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 *LMDBWrapper) NewLMDBTx(write bool, initialIndexName string, frag *fragment) (tx *LMDBTx) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Close shuts down the LMDB database.
|
||||
func (w *LMDBWrapper) Close() (err error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// LMDBTx wraps a lmdb.Txn and provides the Tx interface
|
||||
// method implementations.
|
||||
// The methods on LMDBTx are thread-safe, and can be called
|
||||
// from different goroutines.
|
||||
type LMDBTx struct{}
|
||||
|
||||
func (tx *LMDBTx) Type() string {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (tx *LMDBTx) UseRowCache() bool {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (tx *LMDBTx) Group() *TxGroup {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// 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 *LMDBTx) Pointer() string {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Sn retreives the serial number of the Tx.
|
||||
func (tx *LMDBTx) Sn() int64 {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Rollback rolls back the transaction.
|
||||
func (tx *LMDBTx) Rollback() {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Commit commits the transaction to permanent storage.
|
||||
// Commits can handle up to 100k updates to fragments
|
||||
// at once, but not more. This is a LMDBDB imposed limit.
|
||||
func (tx *LMDBTx) Commit() error {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Readonly returns true iff the LMDBTx is read-only.
|
||||
func (tx *LMDBTx) Readonly() bool {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// RoaringBitmap returns the roaring.Bitmap for all bits in the fragment.
|
||||
func (tx *LMDBTx) RoaringBitmap(index, field, view string, shard uint64) (*roaring.Bitmap, error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Container returns the requested roaring.Container, selected by fragment and ckey
|
||||
func (tx *LMDBTx) Container(index, field, view string, shard uint64, ckey uint64) (c *roaring.Container, err error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// PutContainer stores rc under the specified fragment and container ckey.
|
||||
func (tx *LMDBTx) PutContainer(index, field, view string, shard uint64, ckey uint64, rc *roaring.Container) error {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// RemoveContainer deletes the container specified by the shard and container key ckey
|
||||
func (tx *LMDBTx) RemoveContainer(index, field, view string, shard uint64, ckey uint64) error {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Add sets all the a bits hot in the specified fragment.
|
||||
func (tx *LMDBTx) Add(index, field, view string, shard uint64, batched bool, a ...uint64) (changeCount int, err error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
|
||||
}
|
||||
|
||||
// Remove clears all the specified a bits in the chosen fragment.
|
||||
func (tx *LMDBTx) Remove(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Contains returns exists true iff the bit chosen by key is
|
||||
// hot (set to 1) in specified fragment.
|
||||
func (tx *LMDBTx) Contains(index, field, view string, shard uint64, key uint64) (exists bool, err error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (tx *LMDBTx) SliceOfShards(index, field, view, optionalViewPath string) (sliceOfShards []uint64, err error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// 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.
|
||||
//
|
||||
// LMDBTx notes: We auto-stop at the end of this shard, not going beyond.
|
||||
func (tx *LMDBTx) ContainerIterator(index, field, view string, shard uint64, firstRoaringContainerKey uint64) (citer roaring.ContainerIterator, found bool, err error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
|
||||
}
|
||||
|
||||
// LMDBIterator is the iterator returned from a LMDBTx.ContainerIterator() call.
|
||||
// It implements the roaring.ContainerIterator interface.
|
||||
type LMDBIterator struct{}
|
||||
|
||||
// NewLMDBIterator creates an iterator on tx that will
|
||||
// only return lmdbKeys that start with prefix.
|
||||
func NewLMDBIterator(tx *LMDBTx, prefix []byte) (bi *LMDBIterator) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Close tells the database and transaction that the user is done
|
||||
// with the iterator.
|
||||
func (bi *LMDBIterator) Close() {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Valid returns false if there are no more values in the iterator's range.
|
||||
func (bi *LMDBIterator) Valid() bool {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Seek allows the iterator to start at needle instead of the global begining.
|
||||
func (bi *LMDBIterator) Seek(needle []byte) (ok bool) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (bi *LMDBIterator) ValidForPrefix(prefix []byte) bool {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (bi *LMDBIterator) String() (r string) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
var oneByteSliceOfZero = []byte{0}
|
||||
|
||||
// Next advances the iterator.
|
||||
func (bi *LMDBIterator) Next() (ok bool) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Value retrieves what is pointed at currently by the iterator.
|
||||
func (bi *LMDBIterator) Value() (containerKey uint64, c *roaring.Container) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// lmdbFinder implements roaring.IteratorFinder.
|
||||
// It is used by LMDBTx.ForEach()
|
||||
type lmdbFinder struct {
|
||||
tx *LMDBTx
|
||||
index string
|
||||
field string
|
||||
view string
|
||||
shard uint64
|
||||
needClose []Closer
|
||||
}
|
||||
|
||||
// FindIterator lets lmdbFinder implement the roaring.FindIterator interface.
|
||||
func (bf *lmdbFinder) FindIterator(seek uint64) (roaring.ContainerIterator, bool) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Close closes all bf.needClose listed Closers.
|
||||
func (bf *lmdbFinder) Close() {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// NewTxIterator returns a *roaring.Iterator that MUST have Close() called on it BEFORE
|
||||
// the transaction Commits or Rollsback.
|
||||
func (tx *LMDBTx) NewTxIterator(index, field, view string, shard uint64) *roaring.Iterator {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// ForEach applies fn to each bitmap in the fragment.
|
||||
func (tx *LMDBTx) ForEach(index, field, view string, shard uint64, fn func(i uint64) error) error {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// ForEachRange applies fn on the selected range of bits on the chosen fragment.
|
||||
func (tx *LMDBTx) ForEachRange(index, field, view string, shard uint64, start, end uint64, fn func(uint64) error) error {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Count operates on the full bitmap level, so it sums over all the containers
|
||||
// in the bitmap.
|
||||
func (tx *LMDBTx) Count(index, field, view string, shard uint64) (uint64, error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// 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 *LMDBTx) Max(index, field, view string, shard uint64) (uint64, error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// Min returns the smallest bit set in the fragment. If no bit is hot,
|
||||
// the second return argument is false.
|
||||
func (tx *LMDBTx) Min(index, field, view string, shard uint64) (uint64, bool, error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// 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 *LMDBTx) CountRange(index, field, view string, shard uint64, start, end uint64) (n uint64, err error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// 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 *LMDBTx) OffsetRange(index, field, view string, shard, offset, start, endx uint64) (other *roaring.Bitmap, err error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// IncrementOpN increments the tx opcount by changedN
|
||||
func (tx *LMDBTx) IncrementOpN(index, field, view string, shard uint64, changedN int) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// ImportRoaringBits handles deletes by setting clear=true.
|
||||
// rowSet[rowID] returns the number of bit changed on that rowID.
|
||||
func (tx *LMDBTx) 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) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (tx *LMDBTx) toContainer(typ byte, v []byte) (r *roaring.Container) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// StringifiedLMDBKeys returns a string with all the container
|
||||
// keys available in lmdb.
|
||||
func (w *LMDBWrapper) StringifiedLMDBKeys(optionalUseThisTx Tx) (r string) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// countBitsSet returns the number of bits set (or "hot") in
|
||||
// the roaring container value found by the txkey.Key()
|
||||
// formatted bkey.
|
||||
func (tx *LMDBTx) countBitsSet(bkey []byte) (n int) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (tx *LMDBTx) IsDone() (done bool) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (tx *LMDBTx) Dump(short bool, shard uint64) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// stringifiedLMDBKeysTx reports all the lmdb keys and a
|
||||
// corresponding blake3 hash viewable by txn within the entire
|
||||
// lmdb 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 stringifiedLMDBKeysTx(tx *LMDBTx) (r string) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (w *LMDBWrapper) DeleteField(index, field, fieldPath string) error {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (w *LMDBWrapper) DeleteFragment(index, field, view string, shard uint64, frag interface{}) error {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (w *LMDBWrapper) DeletePrefix(prefix []byte) error {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (tx *LMDBTx) RoaringBitmapReader(index, field, view string, shard uint64, fragmentPathForRoaring string) (r io.ReadCloser, sz int64, err error) {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
// UnionInPlace unions all the others Bitmaps into a new Bitmap, and then writes it to the
|
||||
// specified fragment.
|
||||
func (tx *LMDBTx) UnionInPlace(index, field, view string, shard uint64, others ...*roaring.Bitmap) error {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
|
||||
func (tx *LMDBTx) Options() Txo {
|
||||
panic("lmdb only available on 64-bit arch")
|
||||
}
|
||||
1294
lmdb_test.go
1294
lmdb_test.go
File diff suppressed because it is too large
Load diff
|
|
@ -161,6 +161,7 @@ type RowResponse struct {
|
|||
Headers []*ColumnInfo `protobuf:"bytes,1,rep,name=headers,proto3" json:"headers,omitempty"`
|
||||
Columns []*ColumnResponse `protobuf:"bytes,2,rep,name=columns,proto3" json:"columns,omitempty"`
|
||||
StatusError *StatusError `protobuf:"bytes,3,opt,name=StatusError,proto3" json:"StatusError,omitempty"`
|
||||
Duration int64 `protobuf:"varint,4,opt,name=duration,proto3" json:"duration,omitempty"`
|
||||
XXX_NoUnkeyedLiteral struct{} `json:"-"`
|
||||
XXX_unrecognized []byte `json:"-"`
|
||||
XXX_sizecache int32 `json:"-"`
|
||||
|
|
@ -212,6 +213,13 @@ func (m *RowResponse) GetStatusError() *StatusError {
|
|||
return nil
|
||||
}
|
||||
|
||||
func (m *RowResponse) GetDuration() int64 {
|
||||
if m != nil {
|
||||
return m.Duration
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
type Row struct {
|
||||
Columns []*ColumnResponse `protobuf:"bytes,1,rep,name=columns,proto3" json:"columns,omitempty"`
|
||||
XXX_NoUnkeyedLiteral struct{} `json:"-"`
|
||||
|
|
@ -255,6 +263,7 @@ type TableResponse struct {
|
|||
Headers []*ColumnInfo `protobuf:"bytes,1,rep,name=headers,proto3" json:"headers,omitempty"`
|
||||
Rows []*Row `protobuf:"bytes,2,rep,name=rows,proto3" json:"rows,omitempty"`
|
||||
StatusError *StatusError `protobuf:"bytes,3,opt,name=StatusError,proto3" json:"StatusError,omitempty"`
|
||||
Duration int64 `protobuf:"varint,4,opt,name=duration,proto3" json:"duration,omitempty"`
|
||||
XXX_NoUnkeyedLiteral struct{} `json:"-"`
|
||||
XXX_unrecognized []byte `json:"-"`
|
||||
XXX_sizecache int32 `json:"-"`
|
||||
|
|
@ -306,6 +315,13 @@ func (m *TableResponse) GetStatusError() *StatusError {
|
|||
return nil
|
||||
}
|
||||
|
||||
func (m *TableResponse) GetDuration() int64 {
|
||||
if m != nil {
|
||||
return m.Duration
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
type ColumnInfo struct {
|
||||
Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"`
|
||||
Datatype string `protobuf:"bytes,2,opt,name=datatype,proto3" json:"datatype,omitempty"`
|
||||
|
|
@ -841,54 +857,55 @@ func init() {
|
|||
func init() { proto.RegisterFile("pilosa.proto", fileDescriptor_ef0691a44d1e275c) }
|
||||
|
||||
var fileDescriptor_ef0691a44d1e275c = []byte{
|
||||
// 745 bytes of a gzipped FileDescriptorProto
|
||||
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x9c, 0x55, 0xdd, 0x72, 0xd3, 0x3a,
|
||||
0x10, 0x8e, 0x6b, 0x37, 0x89, 0x37, 0xfd, 0x3b, 0xea, 0x39, 0x3d, 0x99, 0xce, 0x99, 0x83, 0xeb,
|
||||
0x5e, 0x10, 0x06, 0xa6, 0x2d, 0x81, 0xc2, 0x00, 0xe5, 0xa2, 0x2d, 0x30, 0xe9, 0x00, 0x43, 0xaa,
|
||||
0xd2, 0x5e, 0x70, 0xa7, 0xc4, 0x4a, 0xea, 0x41, 0xb1, 0x12, 0xcb, 0x69, 0xc9, 0x0b, 0xf0, 0x06,
|
||||
0xbc, 0x01, 0x6f, 0xc1, 0x3d, 0xcf, 0xc5, 0x48, 0xb2, 0x1c, 0xbb, 0x10, 0xa6, 0xf4, 0xca, 0xda,
|
||||
0xfd, 0xbe, 0xd5, 0xee, 0x6a, 0x7f, 0x0c, 0x0b, 0xc3, 0x90, 0x71, 0x41, 0xb6, 0x86, 0x31, 0x4f,
|
||||
0x38, 0x2a, 0x6b, 0xc9, 0x7f, 0x02, 0xcb, 0xc7, 0x63, 0x1a, 0x4f, 0xda, 0xc7, 0x6f, 0x30, 0x1d,
|
||||
0x8d, 0xa9, 0x48, 0xd0, 0xdf, 0x30, 0x1f, 0x46, 0x01, 0xfd, 0x54, 0xb7, 0x3c, 0xab, 0xe1, 0x62,
|
||||
0x2d, 0xa0, 0x15, 0xb0, 0x87, 0x23, 0x56, 0x9f, 0x53, 0x3a, 0x79, 0xf4, 0x37, 0x53, 0xd3, 0x93,
|
||||
0xa9, 0xe9, 0x0a, 0xd8, 0x62, 0xc4, 0x52, 0x43, 0x79, 0xf4, 0x9f, 0x41, 0xed, 0x24, 0x21, 0xc9,
|
||||
0x58, 0xbc, 0x8c, 0x63, 0x1e, 0x23, 0x04, 0xce, 0x21, 0x0f, 0xa8, 0x62, 0x2c, 0x62, 0x75, 0x46,
|
||||
0x75, 0xa8, 0xbc, 0xa5, 0x42, 0x90, 0x3e, 0x4d, 0x6f, 0x37, 0xa2, 0xff, 0xd5, 0x82, 0x1a, 0xe6,
|
||||
0x97, 0x98, 0x8a, 0x21, 0x8f, 0x04, 0x45, 0xf7, 0xa0, 0x72, 0x4e, 0x49, 0x40, 0x63, 0x51, 0xb7,
|
||||
0x3c, 0xbb, 0x51, 0x6b, 0xa2, 0xad, 0x34, 0xa9, 0x43, 0xce, 0xc6, 0x83, 0xe8, 0x28, 0xea, 0x71,
|
||||
0x6c, 0x28, 0x68, 0x07, 0x2a, 0x5d, 0xa5, 0x16, 0xf5, 0x39, 0xc5, 0x5e, 0x2b, 0xb2, 0xcd, 0xb5,
|
||||
0xd8, 0xd0, 0xd0, 0x6e, 0x21, 0xd8, 0xba, 0xed, 0x59, 0x8d, 0x5a, 0x73, 0xd5, 0x58, 0xe5, 0x20,
|
||||
0x9c, 0xe7, 0xf9, 0x8f, 0xc1, 0xc6, 0xfc, 0x32, 0xef, 0xcf, 0xba, 0x96, 0x3f, 0xff, 0x8b, 0x05,
|
||||
0x8b, 0xef, 0x49, 0x87, 0xd1, 0x1b, 0x66, 0x78, 0x0b, 0x9c, 0x98, 0x5f, 0x9a, 0xf4, 0x6a, 0x86,
|
||||
0x2a, 0x9f, 0x4c, 0x01, 0x37, 0x4d, 0x68, 0x0f, 0x60, 0xea, 0x4e, 0xd6, 0x2c, 0x22, 0x03, 0x9a,
|
||||
0x56, 0x55, 0x9d, 0xd1, 0x3a, 0x54, 0x03, 0x92, 0x90, 0x64, 0x32, 0x34, 0x45, 0xcb, 0x64, 0xff,
|
||||
0xb3, 0x0d, 0x4b, 0xc5, 0x8c, 0xd1, 0xff, 0xe0, 0x8a, 0x24, 0x0e, 0xa3, 0xfe, 0x19, 0x49, 0xbb,
|
||||
0xa3, 0x55, 0xc2, 0x53, 0x95, 0xc4, 0xc7, 0x61, 0x94, 0x3c, 0x7a, 0x28, 0x71, 0x79, 0x9f, 0x23,
|
||||
0xf1, 0x4c, 0x85, 0xfe, 0x83, 0x6a, 0x06, 0xcb, 0x24, 0xec, 0x56, 0x09, 0x67, 0x1a, 0xb4, 0x0e,
|
||||
0x95, 0x0e, 0xe7, 0x4c, 0x82, 0x8e, 0x67, 0x35, 0xaa, 0xad, 0x12, 0x36, 0x0a, 0x85, 0x31, 0xde,
|
||||
0x91, 0xd8, 0xbc, 0x67, 0x35, 0x16, 0x14, 0xa6, 0x15, 0xe8, 0x39, 0x2c, 0x69, 0x17, 0xfb, 0x71,
|
||||
0x4c, 0x26, 0x92, 0x52, 0x2e, 0x3e, 0xd0, 0xe9, 0x14, 0x6d, 0x95, 0xf0, 0x15, 0xb2, 0x34, 0xd7,
|
||||
0x19, 0x64, 0xe6, 0x95, 0xab, 0xef, 0x9b, 0xa1, 0xd2, 0xbc, 0x48, 0x46, 0x1e, 0x40, 0x8f, 0x71,
|
||||
0x92, 0x66, 0x55, 0xf5, 0xac, 0x86, 0xd5, 0x2a, 0xe1, 0x9c, 0x0e, 0xdd, 0x07, 0x08, 0x68, 0x37,
|
||||
0x1c, 0x10, 0x95, 0x9a, 0xab, 0x2e, 0x5f, 0x36, 0x97, 0xbf, 0xd0, 0x88, 0x34, 0x99, 0x92, 0x0e,
|
||||
0x6a, 0xe0, 0xea, 0xe6, 0x3a, 0x23, 0xcc, 0xdf, 0x85, 0x4a, 0xca, 0x92, 0x33, 0x7d, 0x41, 0xd8,
|
||||
0x58, 0x17, 0xd1, 0xc6, 0x5a, 0x90, 0x5a, 0xd1, 0x25, 0x4c, 0x97, 0xd0, 0xc6, 0x5a, 0xf0, 0xbf,
|
||||
0x59, 0xb0, 0x74, 0x14, 0x89, 0x21, 0xed, 0x26, 0xbf, 0x5f, 0x09, 0x77, 0xf3, 0x03, 0x26, 0x83,
|
||||
0xfb, 0xcb, 0x04, 0x77, 0x14, 0x88, 0x77, 0xf1, 0x6b, 0x3a, 0x11, 0xd3, 0xd9, 0xf2, 0x61, 0xa1,
|
||||
0x17, 0xb2, 0x84, 0xc6, 0xaf, 0x42, 0xca, 0x02, 0x51, 0xb7, 0x3d, 0xbb, 0xe1, 0xe2, 0x82, 0x4e,
|
||||
0xba, 0x61, 0xe1, 0x20, 0x4c, 0x54, 0x19, 0x1d, 0xac, 0x05, 0xb4, 0x06, 0x65, 0xde, 0xeb, 0x09,
|
||||
0x9a, 0xa8, 0x0a, 0x3a, 0x38, 0x95, 0x24, 0x7b, 0x24, 0xf7, 0x8f, 0xaa, 0x9a, 0x8b, 0xb5, 0xe0,
|
||||
0x6f, 0x40, 0x2d, 0x57, 0x36, 0xd9, 0xbc, 0x17, 0x84, 0xe9, 0x69, 0x72, 0xb0, 0x3a, 0x4b, 0x4a,
|
||||
0xae, 0x34, 0x05, 0x8a, 0x9b, 0x52, 0xfa, 0xe0, 0x66, 0x39, 0xa0, 0xdb, 0x60, 0x87, 0x81, 0x50,
|
||||
0xb9, 0xcf, 0x6c, 0x0e, 0xc9, 0x40, 0x77, 0xc0, 0xf9, 0x48, 0x27, 0xe6, 0x35, 0x66, 0xf4, 0x81,
|
||||
0xa2, 0x1c, 0x94, 0xc1, 0x91, 0xc3, 0xd2, 0xfc, 0x3e, 0x07, 0xe5, 0xb6, 0xa2, 0xa1, 0x3d, 0xa8,
|
||||
0x9a, 0x7d, 0x8a, 0xfe, 0x35, 0xb6, 0x57, 0x36, 0xec, 0xfa, 0x6a, 0x7e, 0xc8, 0xd3, 0xf1, 0xf2,
|
||||
0x4b, 0x3b, 0x16, 0xda, 0x87, 0x45, 0xc3, 0x3d, 0x8d, 0x48, 0x3c, 0x99, 0x7d, 0xc5, 0x3f, 0x06,
|
||||
0x28, 0xac, 0x1e, 0xbf, 0x94, 0x05, 0xd0, 0xfe, 0x29, 0x80, 0xf6, 0x1f, 0x04, 0xd0, 0xfe, 0x75,
|
||||
0x00, 0xed, 0x6b, 0x04, 0xf0, 0x14, 0x2a, 0x69, 0xe3, 0xa1, 0x6c, 0x77, 0x16, 0x3b, 0x71, 0xa6,
|
||||
0xfb, 0x83, 0xcd, 0x0f, 0x1b, 0xfd, 0x30, 0x39, 0x1f, 0x77, 0xb6, 0xba, 0x7c, 0xb0, 0xad, 0x49,
|
||||
0xe6, 0x73, 0xd1, 0xdc, 0x56, 0x7f, 0xbd, 0x4e, 0x59, 0x7d, 0x1e, 0xfc, 0x08, 0x00, 0x00, 0xff,
|
||||
0xff, 0x60, 0xce, 0x2e, 0x49, 0x0c, 0x07, 0x00, 0x00,
|
||||
// 761 bytes of a gzipped FileDescriptorProto
|
||||
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0xb4, 0x55, 0xcd, 0x72, 0xd3, 0x48,
|
||||
0x10, 0xb6, 0x22, 0xc5, 0xb6, 0xda, 0xf9, 0xdb, 0xc9, 0x6e, 0x56, 0x95, 0xda, 0xda, 0x55, 0x94,
|
||||
0xc3, 0x7a, 0x6b, 0xb7, 0x92, 0xac, 0x77, 0x03, 0x05, 0x84, 0x43, 0x12, 0xa0, 0x9c, 0x02, 0x0a,
|
||||
0x67, 0x42, 0x72, 0xe0, 0x36, 0xb6, 0xc6, 0x8e, 0x8a, 0xb1, 0xc6, 0xd6, 0x48, 0x09, 0x7e, 0x01,
|
||||
0xde, 0x87, 0x33, 0x17, 0x4e, 0x3c, 0x17, 0x35, 0x33, 0x1a, 0xd9, 0x0a, 0x98, 0x0a, 0x54, 0x71,
|
||||
0xf2, 0x74, 0x7f, 0x5f, 0xb7, 0xfa, 0x9b, 0xee, 0x69, 0xc3, 0xd2, 0x28, 0x62, 0x5c, 0x90, 0x9d,
|
||||
0x51, 0xc2, 0x53, 0x8e, 0xaa, 0xda, 0x0a, 0xee, 0xc1, 0xea, 0x69, 0x46, 0x93, 0x49, 0xe7, 0xf4,
|
||||
0x19, 0xa6, 0xe3, 0x8c, 0x8a, 0x14, 0xfd, 0x0c, 0x8b, 0x51, 0x1c, 0xd2, 0x37, 0x9e, 0xe5, 0x5b,
|
||||
0x4d, 0x17, 0x6b, 0x03, 0xad, 0x81, 0x3d, 0x1a, 0x33, 0x6f, 0x41, 0xf9, 0xe4, 0x31, 0xd8, 0xce,
|
||||
0x43, 0xcf, 0xa6, 0xa1, 0x6b, 0x60, 0x8b, 0x31, 0xcb, 0x03, 0xe5, 0x31, 0x78, 0x00, 0x8d, 0xb3,
|
||||
0x94, 0xa4, 0x99, 0x78, 0x9c, 0x24, 0x3c, 0x41, 0x08, 0x9c, 0x63, 0x1e, 0x52, 0xc5, 0x58, 0xc6,
|
||||
0xea, 0x8c, 0x3c, 0xa8, 0x3d, 0xa7, 0x42, 0x90, 0x01, 0xcd, 0xb3, 0x1b, 0x33, 0xf8, 0x60, 0x41,
|
||||
0x03, 0xf3, 0x6b, 0x4c, 0xc5, 0x88, 0xc7, 0x82, 0xa2, 0x7f, 0xa0, 0x76, 0x49, 0x49, 0x48, 0x13,
|
||||
0xe1, 0x59, 0xbe, 0xdd, 0x6c, 0xb4, 0xd0, 0x4e, 0x2e, 0xea, 0x98, 0xb3, 0x6c, 0x18, 0x9f, 0xc4,
|
||||
0x7d, 0x8e, 0x0d, 0x05, 0xed, 0x41, 0xad, 0xa7, 0xdc, 0xc2, 0x5b, 0x50, 0xec, 0x8d, 0x32, 0xdb,
|
||||
0xa4, 0xc5, 0x86, 0x86, 0xf6, 0x4b, 0xc5, 0x7a, 0xb6, 0x6f, 0x35, 0x1b, 0xad, 0x75, 0x13, 0x35,
|
||||
0x03, 0xe1, 0x92, 0xa8, 0x4d, 0xa8, 0x87, 0x59, 0x42, 0xd2, 0x88, 0xc7, 0x9e, 0xe3, 0x5b, 0x4d,
|
||||
0x1b, 0x17, 0x76, 0x70, 0x17, 0x6c, 0xcc, 0xaf, 0x67, 0x6b, 0xb1, 0x6e, 0x55, 0x4b, 0xf0, 0xce,
|
||||
0x82, 0xe5, 0x97, 0xa4, 0xcb, 0xe8, 0x77, 0xaa, 0xff, 0x03, 0x9c, 0x84, 0x5f, 0x1b, 0xe9, 0x0d,
|
||||
0x43, 0x95, 0xd7, 0xa9, 0x80, 0x1f, 0x21, 0xf6, 0x00, 0x60, 0x5a, 0x8a, 0xec, 0x75, 0x4c, 0x86,
|
||||
0x34, 0x9f, 0x06, 0x75, 0x56, 0xd1, 0x24, 0x25, 0xe9, 0x64, 0x64, 0x9a, 0x5d, 0xd8, 0xc1, 0x5b,
|
||||
0x1b, 0x56, 0xca, 0xb7, 0x81, 0x7e, 0x07, 0x57, 0xa4, 0x49, 0x14, 0x0f, 0x2e, 0x48, 0x3e, 0x55,
|
||||
0xed, 0x0a, 0x9e, 0xba, 0x24, 0x9e, 0x45, 0x71, 0x7a, 0xe7, 0x7f, 0x89, 0xcb, 0x7c, 0x8e, 0xc4,
|
||||
0x0b, 0x17, 0xfa, 0x0d, 0xea, 0x05, 0x2c, 0x05, 0xda, 0xed, 0x0a, 0x2e, 0x3c, 0x68, 0x13, 0x6a,
|
||||
0x5d, 0xce, 0x99, 0x04, 0xa5, 0x92, 0x7a, 0xbb, 0x82, 0x8d, 0x43, 0x61, 0x8c, 0x77, 0x25, 0xb6,
|
||||
0xe8, 0x5b, 0xcd, 0x25, 0x85, 0x69, 0x07, 0x7a, 0x08, 0x2b, 0xfa, 0x13, 0x87, 0x49, 0x42, 0x26,
|
||||
0x92, 0x52, 0x2d, 0x5f, 0xde, 0xf9, 0x14, 0x6d, 0x57, 0xf0, 0x0d, 0xb2, 0x0c, 0xd7, 0x0a, 0x8a,
|
||||
0xf0, 0xda, 0xcd, 0xbb, 0x2f, 0x50, 0x19, 0x5e, 0x26, 0x23, 0x1f, 0xa0, 0xcf, 0x38, 0xc9, 0x55,
|
||||
0xd5, 0x7d, 0xab, 0x69, 0xb5, 0x2b, 0x78, 0xc6, 0x87, 0xfe, 0x05, 0x08, 0x69, 0x2f, 0x1a, 0x12,
|
||||
0x25, 0xcd, 0x55, 0xc9, 0x57, 0x4d, 0xf2, 0x47, 0x1a, 0x91, 0x21, 0x53, 0xd2, 0x51, 0x03, 0x5c,
|
||||
0x3d, 0x78, 0x17, 0x84, 0x05, 0xfb, 0x50, 0xcb, 0x59, 0x72, 0x17, 0x5c, 0x11, 0x96, 0xe9, 0x26,
|
||||
0xda, 0x58, 0x1b, 0xd2, 0x2b, 0x7a, 0x84, 0xe9, 0x16, 0xda, 0x58, 0x1b, 0xc1, 0x7b, 0x0b, 0x56,
|
||||
0x4e, 0x62, 0x31, 0xa2, 0xbd, 0xf4, 0xeb, 0xab, 0xe4, 0xef, 0xd9, 0x87, 0x29, 0x8b, 0xfb, 0xc9,
|
||||
0x14, 0x77, 0x12, 0x8a, 0x17, 0xc9, 0x53, 0x3a, 0x11, 0xd3, 0x37, 0x19, 0xc0, 0x52, 0x3f, 0x62,
|
||||
0x29, 0x4d, 0x9e, 0x44, 0x94, 0x85, 0xc2, 0xb3, 0x7d, 0xbb, 0xe9, 0xe2, 0x92, 0x4f, 0x7e, 0x86,
|
||||
0x45, 0xc3, 0x28, 0x55, 0x6d, 0x74, 0xb0, 0x36, 0xd0, 0x06, 0x54, 0x79, 0xbf, 0x2f, 0x68, 0xaa,
|
||||
0x3a, 0xe8, 0xe0, 0xdc, 0x92, 0xec, 0xb1, 0xdc, 0x5b, 0xaa, 0x6b, 0x2e, 0xd6, 0x46, 0xb0, 0x05,
|
||||
0x8d, 0x99, 0xb6, 0xc9, 0xe1, 0xbd, 0x22, 0x4c, 0xbf, 0x34, 0x07, 0xab, 0xb3, 0xa4, 0xcc, 0xb4,
|
||||
0xa6, 0x44, 0x71, 0x73, 0xca, 0x00, 0xdc, 0x42, 0x03, 0xfa, 0x13, 0xec, 0x28, 0x14, 0x4a, 0xfb,
|
||||
0xdc, 0xe1, 0x90, 0x0c, 0xf4, 0x17, 0x38, 0xaf, 0xe9, 0xc4, 0xdc, 0xc6, 0x9c, 0x39, 0x50, 0x94,
|
||||
0xa3, 0x2a, 0x38, 0xf2, 0xb1, 0xb4, 0x3e, 0x2e, 0x40, 0xb5, 0xa3, 0x68, 0xe8, 0x00, 0xea, 0x66,
|
||||
0x0f, 0xa3, 0x5f, 0x4d, 0xec, 0x8d, 0xcd, 0xbc, 0xb9, 0x3e, 0xbb, 0x00, 0xf2, 0xe7, 0x15, 0x54,
|
||||
0xf6, 0x2c, 0x74, 0x08, 0xcb, 0x86, 0x7b, 0x1e, 0x93, 0x64, 0x32, 0x3f, 0xc5, 0x2f, 0x06, 0x28,
|
||||
0xad, 0xa5, 0xa0, 0x52, 0x14, 0xd0, 0xf9, 0xac, 0x80, 0xce, 0x37, 0x14, 0xd0, 0xf9, 0x72, 0x01,
|
||||
0x9d, 0x5b, 0x14, 0x70, 0x1f, 0x6a, 0xf9, 0xe0, 0xa1, 0x62, 0xaf, 0x96, 0x27, 0x71, 0xee, 0xe7,
|
||||
0x8f, 0xb6, 0x5f, 0x6d, 0x0d, 0xa2, 0xf4, 0x32, 0xeb, 0xee, 0xf4, 0xf8, 0x70, 0x57, 0x93, 0xcc,
|
||||
0xcf, 0x55, 0x6b, 0x57, 0xfd, 0x5b, 0x76, 0xab, 0xea, 0xe7, 0xbf, 0x4f, 0x01, 0x00, 0x00, 0xff,
|
||||
0xff, 0x18, 0x84, 0x1c, 0xb4, 0x44, 0x07, 0x00, 0x00,
|
||||
}
|
||||
|
||||
// Reference imports to suppress errors if they are not otherwise used.
|
||||
|
|
|
|||
|
|
@ -23,6 +23,7 @@ message RowResponse{
|
|||
repeated ColumnInfo headers = 1;
|
||||
repeated ColumnResponse columns = 2;
|
||||
StatusError StatusError = 3;
|
||||
int64 duration = 4;
|
||||
}
|
||||
|
||||
message Row {
|
||||
|
|
@ -33,6 +34,7 @@ message TableResponse{
|
|||
repeated ColumnInfo headers = 1;
|
||||
repeated Row rows = 2;
|
||||
StatusError StatusError = 3;
|
||||
int64 duration = 4;
|
||||
}
|
||||
|
||||
message ColumnInfo {
|
||||
|
|
|
|||
191
rbf/cursor.go
191
rbf/cursor.go
|
|
@ -33,7 +33,6 @@ type Cursor struct {
|
|||
buffered bool
|
||||
|
||||
// buffers
|
||||
leafPage []byte
|
||||
array [ArrayMaxSize + 1]uint16
|
||||
rle [RLEMaxSize + 1]roaring.Interval16
|
||||
leafCells [PageSize / 8]leafCell
|
||||
|
|
@ -140,7 +139,13 @@ func (c *Cursor) Add(v uint64) (changed bool, err error) {
|
|||
}
|
||||
|
||||
// If the container exists and bit is not set then update the page.
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
switch cell.Type {
|
||||
case ContainerTypeArray:
|
||||
// Exit if value exists in array container.
|
||||
|
|
@ -205,7 +210,13 @@ func (c *Cursor) Remove(v uint64) (changed bool, err error) {
|
|||
}
|
||||
|
||||
// If the container exists and bit is not set then update the page.
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
switch cell.Type {
|
||||
case ContainerTypeArray:
|
||||
// Exit if value does not exists in array container.
|
||||
|
|
@ -315,7 +326,13 @@ func (c *Cursor) Contains(v uint64) (exists bool, err error) {
|
|||
}
|
||||
|
||||
// If the container exists then check for low bits existence.
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
switch cell.Type {
|
||||
case ContainerTypeArray:
|
||||
a := toArray16(cell.Data)
|
||||
|
|
@ -351,14 +368,16 @@ func toPgno(val []byte) uint32 {
|
|||
}
|
||||
|
||||
func (c *Cursor) putLeafCell(in leafCell) (err error) {
|
||||
|
||||
leafPage := c.leafPage // the last read leaf page
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, isHeap, err := c.tx.readPage(elem.pgno) // the last read leaf page
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
cellN := readCellN(leafPage)
|
||||
|
||||
// Determine if the insert/update will overflow the page.
|
||||
// If it doesn't then we can do an optimized write where we don't deserialize.
|
||||
isInsert := elem.index >= cellN || c.Key() != in.Key
|
||||
isInsert := elem.index >= cellN || pageKeyAt(leafPage, elem.index) != in.Key
|
||||
newEstPageSize := leafPageSize(leafPage)
|
||||
if isInsert {
|
||||
newEstPageSize += in.Size() + leafCellIndexElemSize
|
||||
|
|
@ -477,6 +496,11 @@ func (c *Cursor) putLeafCell(in leafCell) (err error) {
|
|||
parents = append(parents, parent)
|
||||
}
|
||||
|
||||
// Free the source page once we've finished with it if it is on heap.
|
||||
if isHeap {
|
||||
freePage(leafPage)
|
||||
}
|
||||
|
||||
// TODO(BBJ): Update page in buffer & cursor stack.
|
||||
|
||||
// If this is not a split then exit now.
|
||||
|
|
@ -498,8 +522,11 @@ func (c *Cursor) putLeafCell(in leafCell) (err error) {
|
|||
// putLeafCellFast quickly insert or updates a cell on a leaf page.
|
||||
// It works by shifting bytes around instead of deserializing. This must not overflow.
|
||||
func (c *Cursor) putLeafCellFast(in leafCell, isInsert bool) (err error) {
|
||||
src := c.leafPage
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
src, isHeap, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
srcCellN := readCellN(src)
|
||||
|
||||
// Determine the cell count of the new page.
|
||||
|
|
@ -509,34 +536,50 @@ func (c *Cursor) putLeafCellFast(in leafCell, isInsert bool) (err error) {
|
|||
}
|
||||
|
||||
// Write page header.
|
||||
dst := make([]byte, PageSize)
|
||||
dst := allocPage() // make([]byte, PageSize)
|
||||
writePageNo(dst, readPageNo(src))
|
||||
writeFlags(dst, PageTypeLeaf)
|
||||
writeCellN(dst, dstCellN)
|
||||
|
||||
// Loop over source page elements and copy them to the new page.
|
||||
// Copy data before index.
|
||||
offset := dataOffset(dstCellN)
|
||||
for i, j := 0, 0; j < dstCellN; i, j = i+1, j+1 {
|
||||
// If positioned at the insert/update index, write the new cell.
|
||||
if i == elem.index {
|
||||
writeLeafCell(dst[:], j, offset, in)
|
||||
offset += align8(in.Size())
|
||||
|
||||
// If this is an update, skip to the next element.
|
||||
if !isInsert {
|
||||
continue
|
||||
}
|
||||
if elem.index > 0 {
|
||||
srcStart := dataOffset(srcCellN)
|
||||
srcEnd := readCellEndingOffset(src, elem.index-1)
|
||||
shiftN := offset - srcStart
|
||||
copy(dst[offset:], src[srcStart:srcEnd])
|
||||
offset += align8(srcEnd - srcStart)
|
||||
|
||||
// If this is an insert, move the dst position forward.
|
||||
j++
|
||||
// Rewrite initial index slots by position moved.
|
||||
for i := 0; i < elem.index; i++ {
|
||||
writeCellOffset(dst, i, readCellOffset(src, i)+shiftN)
|
||||
}
|
||||
}
|
||||
|
||||
// Copy the raw bytes from the src page to the dst page.
|
||||
if i < srcCellN {
|
||||
srcCellBuf := readLeafCellBytesAtOffset(src, readCellOffset(src, i))
|
||||
writeCellOffset(dst, j, offset)
|
||||
copy(dst[offset:], srcCellBuf)
|
||||
offset += align8(len(srcCellBuf))
|
||||
// Insert new row.
|
||||
writeLeafCell(dst[:], elem.index, offset, in)
|
||||
offset += align8(in.Size())
|
||||
|
||||
// Copy data after inserted element.
|
||||
if (isInsert && elem.index < srcCellN) || (!isInsert && elem.index < srcCellN-1) {
|
||||
var srcStart int
|
||||
if isInsert {
|
||||
srcStart = readCellOffset(src, elem.index)
|
||||
} else {
|
||||
srcStart = readCellOffset(src, elem.index+1)
|
||||
}
|
||||
srcEnd := readCellEndingOffset(src, srcCellN-1)
|
||||
copy(dst[offset:], src[srcStart:srcEnd])
|
||||
|
||||
// Rewrite ending index slots by position moved.
|
||||
shiftN := offset - srcStart
|
||||
for i := elem.index + 1; i < dstCellN; i++ {
|
||||
srci := i
|
||||
if isInsert {
|
||||
srci--
|
||||
}
|
||||
writeCellOffset(dst, i, readCellOffset(src, srci)+shiftN)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -545,15 +588,24 @@ func (c *Cursor) putLeafCellFast(in leafCell, isInsert bool) (err error) {
|
|||
return err
|
||||
}
|
||||
|
||||
// Free page if on heap.
|
||||
if isHeap {
|
||||
freePage(src)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// deleteLeafCell removes a cell from the currently positioned page & index.
|
||||
func (c *Cursor) deleteLeafCell(key uint64) (err error) {
|
||||
cells := readLeafCells(c.leafPage, c.leafCells[:])
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
cells := readLeafCells(leafPage, c.leafCells[:])
|
||||
oldPageKey := cells[0].Key
|
||||
cell := c.cell()
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
if cell.Type == ContainerTypeBitmapPtr {
|
||||
if err := c.tx.freePgno(toPgno(cell.Data)); err != nil {
|
||||
return err
|
||||
|
|
@ -600,7 +652,7 @@ func (c *Cursor) putBranchCells(stackIndex int, newCells []branchCell) (err erro
|
|||
elem := &c.stack.elems[stackIndex]
|
||||
|
||||
// Read branch page from disk. The current buffer is the leaf page.
|
||||
page, err := c.tx.readPage(elem.pgno)
|
||||
page, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -680,7 +732,7 @@ func (c *Cursor) updateBranchCell(stackIndex int, newKey uint64) (err error) {
|
|||
elem := &c.stack.elems[stackIndex]
|
||||
|
||||
// Read branch page from disk. The current buffer is the leaf page.
|
||||
page, err := c.tx.readPage(elem.pgno)
|
||||
page, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -716,7 +768,7 @@ func (c *Cursor) deleteBranchCell(stackIndex int, key uint64) (err error) {
|
|||
elem := &c.stack.elems[stackIndex]
|
||||
|
||||
// Read branch page from disk. The current buffer is the leaf page.
|
||||
page, err := c.tx.readPage(elem.pgno)
|
||||
page, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -730,7 +782,7 @@ func (c *Cursor) deleteBranchCell(stackIndex int, key uint64) (err error) {
|
|||
|
||||
// If the root only has one node, replace it with its child.
|
||||
if stackIndex == 0 && len(cells) == 1 {
|
||||
target, err := c.tx.readPage(cells[0].Pgno)
|
||||
target, _, err := c.tx.readPage(cells[0].Pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -837,16 +889,10 @@ func splitBranchCells(cells []branchCell) [][]branchCell {
|
|||
return slices
|
||||
}
|
||||
|
||||
// Key returns the key that the cursor is currently positioned over.
|
||||
func (c *Cursor) Key() uint64 {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
offset := readCellOffset(c.leafPage, elem.index)
|
||||
return *(*uint64)(unsafe.Pointer(&c.leafPage[offset]))
|
||||
}
|
||||
|
||||
func (c *Cursor) cell() leafCell {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
return readLeafCell(c.leafPage[:], elem.index)
|
||||
// pageKeyAt returns the key at the given index of the page.
|
||||
func pageKeyAt(page []byte, index int) uint64 {
|
||||
offset := readCellOffset(page, index)
|
||||
return *(*uint64)(unsafe.Pointer(&page[offset]))
|
||||
}
|
||||
|
||||
// First moves to the first element of the btree.
|
||||
|
|
@ -856,7 +902,7 @@ func (c *Cursor) First() error {
|
|||
for c.stack.index = 0; ; c.stack.index++ {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
|
||||
buf, err := c.tx.readPage(elem.pgno)
|
||||
buf, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -874,7 +920,6 @@ func (c *Cursor) First() error {
|
|||
}
|
||||
|
||||
case PageTypeLeaf:
|
||||
c.leafPage = buf
|
||||
elem.index = 0
|
||||
if readCellN(buf) == 0 {
|
||||
return io.EOF // root leaf with no elements
|
||||
|
|
@ -894,7 +939,7 @@ func (c *Cursor) Last() error {
|
|||
for c.stack.index = 0; ; c.stack.index++ {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
|
||||
buf, err := c.tx.readPage(elem.pgno)
|
||||
buf, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -912,7 +957,6 @@ func (c *Cursor) Last() error {
|
|||
|
||||
case PageTypeLeaf:
|
||||
elem.index = readCellN(buf) - 1
|
||||
c.leafPage = buf
|
||||
if readCellN(buf) == 0 {
|
||||
return io.EOF // root leaf with no elements
|
||||
}
|
||||
|
|
@ -932,7 +976,7 @@ func (c *Cursor) Seek(key uint64) (exact bool, err error) {
|
|||
elem := &c.stack.elems[c.stack.index]
|
||||
assert(elem.pgno != 0) // cursor should never point to page zero (meta)
|
||||
|
||||
buf, err := c.tx.readPage(elem.pgno)
|
||||
buf, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
|
@ -974,7 +1018,6 @@ func (c *Cursor) Seek(key uint64) (exact bool, err error) {
|
|||
return 1
|
||||
})
|
||||
elem.index = index
|
||||
c.leafPage = buf
|
||||
return xact, nil
|
||||
|
||||
default:
|
||||
|
|
@ -985,18 +1028,24 @@ func (c *Cursor) Seek(key uint64) (exact bool, err error) {
|
|||
|
||||
// Next moves to the next element of the btree. Returns EOF if no more elements exist.
|
||||
func (c *Cursor) Next() error {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if c.buffered {
|
||||
c.buffered = false
|
||||
|
||||
// Move to next available element if we are past the last cell in the page.
|
||||
if elem := &c.stack.elems[c.stack.index]; elem.index >= readCellN(c.leafPage) {
|
||||
if elem.index >= readCellN(leafPage) {
|
||||
return c.goNextPage()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Move forward to the next leaf element if available.
|
||||
if elem := &c.stack.elems[c.stack.index]; elem.index < readCellN(c.leafPage)-1 {
|
||||
if elem.index < readCellN(leafPage)-1 {
|
||||
elem.index++
|
||||
return nil
|
||||
}
|
||||
|
|
@ -1035,7 +1084,7 @@ func (c *Cursor) Prev() error {
|
|||
for ; ; c.stack.index++ {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
|
||||
buf, err := c.tx.readPage(elem.pgno)
|
||||
buf, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -1051,7 +1100,6 @@ func (c *Cursor) Prev() error {
|
|||
|
||||
case PageTypeLeaf:
|
||||
elem.index = readCellN(buf) - 1
|
||||
c.leafPage = buf
|
||||
return nil
|
||||
default:
|
||||
return fmt.Errorf("rbf.Cursor.Prev(): invalid page type: pgno=%d type=%d", elem.pgno, typ)
|
||||
|
|
@ -1059,13 +1107,25 @@ func (c *Cursor) Prev() error {
|
|||
}
|
||||
}
|
||||
|
||||
// Key returns the key for the container the cursor is currently pointing to.
|
||||
func (c *Cursor) Key() uint64 {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, _ := c.tx.readPage(elem.pgno)
|
||||
if readCellN(leafPage[:]) == 0 {
|
||||
return 0
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
return cell.Key
|
||||
}
|
||||
|
||||
// Values returns the values for the container the cursor is currently pointing to.
|
||||
func (c *Cursor) Values() []uint16 {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
if readCellN(c.leafPage[:]) == 0 {
|
||||
leafPage, _, _ := c.tx.readPage(elem.pgno)
|
||||
if readCellN(leafPage[:]) == 0 {
|
||||
return nil
|
||||
}
|
||||
cell := readLeafCell(c.leafPage[:], elem.index)
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
return cell.Values(c.tx)
|
||||
}
|
||||
|
||||
|
|
@ -1079,7 +1139,7 @@ type stackElem struct {
|
|||
func (c *Cursor) goNextPage() error {
|
||||
for c.stack.index--; c.stack.index >= 0; c.stack.index-- {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
if buf, err := c.tx.readPage(elem.pgno); err != nil {
|
||||
if buf, _, err := c.tx.readPage(elem.pgno); err != nil {
|
||||
return err
|
||||
} else if n := readCellN(buf); elem.index+1 < n {
|
||||
elem.index++
|
||||
|
|
@ -1096,7 +1156,7 @@ func (c *Cursor) goNextPage() error {
|
|||
// Traverse back down the stack to find the first element in each page.
|
||||
for ; ; c.stack.index++ {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
buf, err := c.tx.readPage(elem.pgno)
|
||||
buf, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -1110,7 +1170,6 @@ func (c *Cursor) goNextPage() error {
|
|||
}
|
||||
case PageTypeLeaf:
|
||||
elem.index = 0
|
||||
c.leafPage = buf
|
||||
return nil
|
||||
default:
|
||||
return fmt.Errorf("rbf.Cursor.Next(): invalid page type: pgno=%d type=%d", elem.pgno, typ)
|
||||
|
|
@ -1163,7 +1222,13 @@ func ConvertToLeafArgs(key uint64, c *roaring.Container) (result leafCell) {
|
|||
}
|
||||
|
||||
func (c *Cursor) merge(key uint64, data *roaring.Container) (bool, error) {
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
var container *roaring.Container
|
||||
switch cell.Type {
|
||||
case ContainerTypeArray:
|
||||
|
|
@ -1246,7 +1311,13 @@ func (c *Cursor) RemoveRoaring(bm *roaring.Bitmap) (changed bool, err error) {
|
|||
}
|
||||
|
||||
func (c *Cursor) difference(key uint64, data *roaring.Container) (bool, error) {
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
var container *roaring.Container
|
||||
switch cell.Type {
|
||||
case ContainerTypeArray:
|
||||
|
|
|
|||
|
|
@ -61,7 +61,14 @@ func (c *Cursor) Rows() ([]uint64, error) {
|
|||
if err != nil {
|
||||
break
|
||||
}
|
||||
cell := c.cell()
|
||||
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
vRow := cell.Key >> shardVsContainerExponent
|
||||
if vRow == lastRow {
|
||||
continue
|
||||
|
|
@ -92,7 +99,12 @@ func (c *Cursor) DumpKeys() {
|
|||
if err == io.EOF {
|
||||
break
|
||||
}
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
fmt.Println("key", cell.Key)
|
||||
}
|
||||
}
|
||||
|
|
@ -129,7 +141,11 @@ func (c *Cursor) Row(shard, rowID uint64) (*roaring.Bitmap, error) {
|
|||
}
|
||||
if !ok {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
n := readCellN(c.leafPage)
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
n := readCellN(leafPage)
|
||||
if elem.index >= n {
|
||||
if err := c.goNextPage(); err != nil {
|
||||
return nil, errors.Wrap(err, "row")
|
||||
|
|
@ -145,7 +161,12 @@ func (c *Cursor) Row(shard, rowID uint64) (*roaring.Bitmap, error) {
|
|||
return nil, err
|
||||
}
|
||||
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
if cell.Key >= hi1 {
|
||||
break
|
||||
}
|
||||
|
|
@ -157,7 +178,9 @@ func (c *Cursor) Row(shard, rowID uint64) (*roaring.Bitmap, error) {
|
|||
// CurrentPageType returns the type of the container currently pointed to by cursor used in testing
|
||||
// sometimes the cursor needs to be positions prior to this call with First/Last etc.
|
||||
func (c *Cursor) CurrentPageType() ContainerType {
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, _ := c.tx.readPage(elem.pgno)
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
return cell.Type
|
||||
}
|
||||
|
||||
|
|
@ -219,7 +242,7 @@ type Walker interface {
|
|||
}
|
||||
|
||||
func WalkPage(tx *Tx, pgno uint32, walker Walker) {
|
||||
page, err := tx.readPage(pgno)
|
||||
page, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
|
|
|||
18
rbf/db.go
18
rbf/db.go
|
|
@ -599,3 +599,21 @@ func (db *DB) getCursor(tx *Tx) (c *Cursor) {
|
|||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// Shared pool for in-memory database pages.
|
||||
// These are used before being flushed to disk.
|
||||
var pagePool = &sync.Pool{
|
||||
New: func() interface{} {
|
||||
page := make([]byte, PageSize)
|
||||
return &page
|
||||
},
|
||||
}
|
||||
|
||||
func allocPage() []byte {
|
||||
page := pagePool.Get().(*[]byte)
|
||||
return *page
|
||||
}
|
||||
|
||||
func freePage(page []byte) {
|
||||
pagePool.Put(&page)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -62,7 +62,7 @@ func dotCell(b []byte, parent string, writer io.Writer) {
|
|||
|
||||
// dumpdot recursively writes the tree representation starting from a given page to STDERR.
|
||||
func dumpdot(tx *Tx, pgno uint32, parent string, writer io.Writer) {
|
||||
page, err := tx.readPage(pgno)
|
||||
page, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -211,6 +211,12 @@ func writeCellOffset(page []byte, i int, v int) {
|
|||
binary.BigEndian.PutUint16(page[10+(i*2):], uint16(v))
|
||||
}
|
||||
|
||||
// readCellEndingOffset returns the last byte position of the i-th cell.
|
||||
func readCellEndingOffset(page []byte, i int) int {
|
||||
offset := readCellOffset(page, i)
|
||||
return offset + len(readLeafCellBytesAtOffset(page, offset))
|
||||
}
|
||||
|
||||
func dataOffset(n int) int {
|
||||
return align8(10 + (n * 2))
|
||||
}
|
||||
|
|
@ -676,7 +682,7 @@ func Pagedump(b []byte, indent string, writer io.Writer) {
|
|||
|
||||
func Walk(tx *Tx, pgno uint32, v func(uint32, []*RootRecord)) {
|
||||
for pgno := readMetaRootRecordPageNo(tx.meta[:]); pgno != 0; {
|
||||
page, err := tx.readPage(pgno)
|
||||
page, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
|
|
|||
141
rbf/tx.go
141
rbf/tx.go
|
|
@ -373,7 +373,7 @@ func (tx *Tx) RootRecords() (records *immutable.SortedMap, err error) {
|
|||
|
||||
records = immutable.NewSortedMap(nil)
|
||||
for pgno := readMetaRootRecordPageNo(tx.meta[:]); pgno != 0; {
|
||||
page, err := tx.readPage(pgno)
|
||||
page, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -401,7 +401,7 @@ func (tx *Tx) writeRootRecordPages(records *immutable.SortedMap) (err error) {
|
|||
|
||||
// Release all existing root record pages.
|
||||
for pgno := readMetaRootRecordPageNo(tx.meta[:]); pgno != 0; {
|
||||
page, err := tx.readPage(pgno)
|
||||
page, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -598,7 +598,12 @@ func (tx *Tx) RoaringBitmap(name string) (*roaring.Bitmap, error) {
|
|||
return nil, err
|
||||
}
|
||||
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
other.Containers.Put(cell.Key, toContainer(cell, tx))
|
||||
}
|
||||
}
|
||||
|
|
@ -629,7 +634,14 @@ func (tx *Tx) container(name string, key uint64) (*roaring.Container, error) {
|
|||
return nil, err
|
||||
}
|
||||
|
||||
return toContainer(c.cell(), tx), nil
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
return toContainer(cell, tx), nil
|
||||
}
|
||||
|
||||
// PutContainer inserts a container into a bitmap. Overwrites if key already exists.
|
||||
|
|
@ -733,7 +745,7 @@ func (tx *Tx) checkPageAllocations() error {
|
|||
if isInuse && isFree {
|
||||
return fmt.Errorf("page in-use & free: pgno=%d", pgno)
|
||||
} else if !isInuse && !isFree {
|
||||
page, err := tx.readPage(pgno)
|
||||
page, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -767,7 +779,13 @@ func (tx *Tx) freePageSet() (map[uint32]struct{}, error) {
|
|||
return m, err
|
||||
}
|
||||
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
for _, v := range cell.Values(tx) {
|
||||
pgno := uint32((cell.Key << 16) & uint64(v))
|
||||
m[pgno] = struct{}{}
|
||||
|
|
@ -784,7 +802,7 @@ func (tx *Tx) inusePageSet() (map[uint32]struct{}, error) {
|
|||
for pgno := readMetaRootRecordPageNo(tx.meta[:]); pgno != 0; {
|
||||
m[pgno] = struct{}{}
|
||||
|
||||
page, err := tx.readPage(pgno)
|
||||
page, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -822,7 +840,7 @@ func (tx *Tx) inusePageSet() (map[uint32]struct{}, error) {
|
|||
// walkTree recursively iterates over a page and all its children.
|
||||
func (tx *Tx) walkTree(pgno, parent uint32, fn func(pgno, parent, typ uint32) error) error {
|
||||
// Read page and iterate over children.
|
||||
page, err := tx.readPage(pgno)
|
||||
page, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -892,7 +910,13 @@ func (tx *Tx) nextFreelistPageNo() (uint32, error) {
|
|||
return 0, err
|
||||
}
|
||||
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
v := cell.firstValue(tx)
|
||||
|
||||
pgno := uint32((cell.Key << 16) | uint64(v))
|
||||
|
|
@ -914,7 +938,7 @@ func (tx *Tx) freePgno(pgno uint32) error {
|
|||
|
||||
// deallocateTree recursively all pages in a btree.
|
||||
func (tx *Tx) deallocateTree(pgno uint32) error {
|
||||
page, err := tx.readPage(pgno)
|
||||
page, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -936,34 +960,36 @@ func (tx *Tx) deallocateTree(pgno uint32) error {
|
|||
}
|
||||
}
|
||||
|
||||
func (tx *Tx) readPage(pgno uint32) ([]byte, error) {
|
||||
func (tx *Tx) readPage(pgno uint32) (_ []byte, isHeap bool, err error) {
|
||||
// Meta page is always cached on the transaction.
|
||||
if pgno == 0 {
|
||||
return tx.meta[:], nil
|
||||
return tx.meta[:], false, nil
|
||||
}
|
||||
|
||||
// Verify page number requested is within current size of database.
|
||||
pageN := readMetaPageN(tx.meta[:])
|
||||
if pgno > pageN {
|
||||
return nil, fmt.Errorf("rbf: page read out of bounds: pgno=%d max=%d", pgno, pageN)
|
||||
return nil, false, fmt.Errorf("rbf: page read out of bounds: pgno=%d max=%d", pgno, pageN)
|
||||
}
|
||||
|
||||
// Check if page has been updated in this tx.
|
||||
if tx.writable {
|
||||
if page := tx.dirtyPages[pgno]; page != nil {
|
||||
return page, nil
|
||||
return page, true, nil
|
||||
} else if page := tx.dirtyBitmapPages[pgno]; page != nil {
|
||||
return page, nil
|
||||
return page, true, nil
|
||||
}
|
||||
}
|
||||
|
||||
// Check if page is remapped in WAL.
|
||||
if walID, ok := tx.pageMap.Get(pgno); ok {
|
||||
return tx.db.readWALPageByID(walID)
|
||||
buf, err := tx.db.readWALPageByID(walID)
|
||||
return buf, false, err
|
||||
}
|
||||
|
||||
// Otherwise read directly from DB.
|
||||
return tx.db.readDBPage(pgno)
|
||||
buf, err := tx.db.readDBPage(pgno)
|
||||
return buf, false, err
|
||||
}
|
||||
|
||||
func (tx *Tx) writePage(page []byte) error {
|
||||
|
|
@ -1001,7 +1027,7 @@ func (tx *Tx) AddRoaring(name string, bm *roaring.Bitmap) (changed bool, err err
|
|||
}
|
||||
|
||||
func (tx *Tx) leafCellBitmap(pgno uint32) (uint32, []uint64, error) {
|
||||
page, err := tx.readPage(pgno)
|
||||
page, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
return 0, nil, err
|
||||
}
|
||||
|
|
@ -1053,7 +1079,14 @@ func (tx *Tx) ForEachRange(name string, start, end uint64, fn func(uint64) error
|
|||
return err
|
||||
}
|
||||
|
||||
switch cell := c.cell(); cell.Type {
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
switch cell.Type {
|
||||
case ContainerTypeArray:
|
||||
for _, lo := range toArray16(cell.Data) {
|
||||
v := cell.Key<<16 | uint64(lo)
|
||||
|
|
@ -1146,7 +1179,14 @@ func (tx *Tx) Count(name string) (uint64, error) {
|
|||
return 0, err
|
||||
}
|
||||
|
||||
n += uint64(c.cell().BitN)
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
n += uint64(cell.BitN)
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
|
|
@ -1169,7 +1209,13 @@ func (tx *Tx) Max(name string) (uint64, error) {
|
|||
return 0, err
|
||||
}
|
||||
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
return uint64((cell.Key << 16) | uint64(cell.lastValue(tx))), nil
|
||||
}
|
||||
|
||||
|
|
@ -1191,7 +1237,13 @@ func (tx *Tx) Min(name string) (uint64, bool, error) {
|
|||
return 0, false, err
|
||||
}
|
||||
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return 0, false, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
return uint64((cell.Key << 16) | uint64(cell.firstValue(tx))), true, nil
|
||||
}
|
||||
|
||||
|
|
@ -1252,7 +1304,13 @@ func (tx *Tx) CountRange(name string, start, end uint64) (uint64, error) {
|
|||
return 0, err
|
||||
}
|
||||
|
||||
c := csr.cell()
|
||||
elem := &csr.stack.elems[csr.stack.index]
|
||||
leafPage, _, err := csr.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
c := readLeafCell(leafPage, elem.index)
|
||||
|
||||
k := c.Key
|
||||
if k > ekey {
|
||||
break
|
||||
|
|
@ -1324,7 +1382,12 @@ func (tx *Tx) OffsetRange(name string, offset, start, endx uint64) (*roaring.Bit
|
|||
return nil, err
|
||||
}
|
||||
|
||||
cell := c.cell()
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
ckey := cell.Key
|
||||
|
||||
// >= hi1 is correct b/c endx cannot have any lowbits set.
|
||||
|
|
@ -1356,7 +1419,9 @@ func (itr *containerIterator) Next() bool {
|
|||
|
||||
// Value returns the current key & container.
|
||||
func (itr *containerIterator) Value() (uint64, *roaring.Container) {
|
||||
cell := itr.cursor.cell()
|
||||
elem := &itr.cursor.stack.elems[itr.cursor.stack.index]
|
||||
leafPage, _, _ := itr.cursor.tx.readPage(elem.pgno)
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
return cell.Key, toContainer(cell, itr.cursor.tx)
|
||||
}
|
||||
|
||||
|
|
@ -1404,7 +1469,12 @@ func (tx *Tx) DumpString(short bool, shard uint64) (r string) {
|
|||
break
|
||||
}
|
||||
panicOn(err)
|
||||
cell := c.cell()
|
||||
|
||||
elem := &c.stack.elems[c.stack.index]
|
||||
leafPage, _, err := c.tx.readPage(elem.pgno)
|
||||
panicOn(err)
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
|
||||
ckey := cell.Key
|
||||
ct := toContainer(cell, tx)
|
||||
|
||||
|
|
@ -1532,7 +1602,13 @@ func (tx *Tx) ImportRoaringBits(name string, itr roaring.RoaringIterator, clear
|
|||
if exact, err := cur.Seek(itrKey); err != nil {
|
||||
return changed, rowSet, err
|
||||
} else if exact {
|
||||
oldC = toContainer(cur.cell(), tx)
|
||||
elem := &cur.stack.elems[cur.stack.index]
|
||||
leafPage, _, err := cur.tx.readPage(elem.pgno)
|
||||
if err != nil {
|
||||
return changed, rowSet, err
|
||||
}
|
||||
cell := readLeafCell(leafPage, elem.index)
|
||||
oldC = toContainer(cell, tx)
|
||||
}
|
||||
|
||||
if oldC == nil || oldC.N() == 0 {
|
||||
|
|
@ -1687,7 +1763,7 @@ func (tx *Tx) Pages(pgnos []uint32) ([]Page, error) {
|
|||
// Loop over each requested page number and extract additional data.
|
||||
var pages []Page
|
||||
for _, pgno := range pgnos {
|
||||
buf, err := tx.readPage(pgno)
|
||||
buf, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -1805,7 +1881,7 @@ func (tx *Tx) PageInfos() ([]PageInfo, error) {
|
|||
|
||||
// metaPageInfo returns page metadata for the meta page.
|
||||
func (tx *Tx) metaPageInfo() (*MetaPageInfo, error) {
|
||||
buf, err := tx.readPage(0)
|
||||
buf, _, err := tx.readPage(0)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -1822,7 +1898,7 @@ func (tx *Tx) metaPageInfo() (*MetaPageInfo, error) {
|
|||
|
||||
// rootRecordPageInfo returns page metadata for a root record page.
|
||||
func (tx *Tx) rootRecordPageInfo(pgno uint32) (*RootRecordPageInfo, error) {
|
||||
buf, err := tx.readPage(pgno)
|
||||
buf, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -1835,7 +1911,7 @@ func (tx *Tx) rootRecordPageInfo(pgno uint32) (*RootRecordPageInfo, error) {
|
|||
|
||||
func (tx *Tx) walkPageInfo(infos []PageInfo, root uint32, name string) error {
|
||||
return tx.walkTree(root, 0, func(pgno, parent, typ uint32) error {
|
||||
buf, err := tx.readPage(pgno)
|
||||
buf, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -1873,7 +1949,8 @@ func (tx *Tx) walkPageInfo(infos []PageInfo, root uint32, name string) error {
|
|||
|
||||
// PageData returns the raw page data for a single page.
|
||||
func (tx *Tx) PageData(pgno uint32) ([]byte, error) {
|
||||
return tx.readPage(pgno)
|
||||
buf, _, err := tx.readPage(pgno)
|
||||
return buf, err
|
||||
}
|
||||
|
||||
type PageInfo interface {
|
||||
|
|
|
|||
|
|
@ -46,7 +46,7 @@ func (c_orig *Cursor) DebugSlowCheckAllPages() {
|
|||
|
||||
// checkElemNBitN recursively writes the tree representation starting from a given page to STDERR.
|
||||
func checkElemNBitN(tx *Tx, pgno uint32) {
|
||||
page, err := tx.readPage(pgno)
|
||||
page, _, err := tx.readPage(pgno)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1756,6 +1756,14 @@ type RoaringIterator interface {
|
|||
// Clone copies the iterator, preserving it at this point in the iteration.
|
||||
// It may well share much underlying data.
|
||||
Clone() RoaringIterator
|
||||
|
||||
// ContainerKeySpan provides the smallest and largest
|
||||
// container keys that the iterator will return.
|
||||
// The current implementation requires that the underlying header
|
||||
// lists the keys in ascending order.
|
||||
// Iff there no keys, then empty will be returned true.
|
||||
// If there is only a single key, then ckeyLast will equal ckeyFirst.
|
||||
ContainerKeySpan() (ckeyFirst, ckeyLast uint64, empty bool)
|
||||
}
|
||||
|
||||
// baseRoaringIterator holds values used by both Pilosa and official Roaring
|
||||
|
|
@ -1924,6 +1932,22 @@ func (r *baseRoaringIterator) Done(err error) {
|
|||
r.currentDataOffset = 0
|
||||
}
|
||||
|
||||
func (r *baseRoaringIterator) ContainerKeySpan() (ckeyFirst, ckeyLast uint64, empty bool) {
|
||||
n := r.keys
|
||||
if n == 0 {
|
||||
empty = true
|
||||
return
|
||||
}
|
||||
ckeyFirst = binary.LittleEndian.Uint64(r.headers[0:8])
|
||||
if n == 1 {
|
||||
ckeyLast = ckeyFirst
|
||||
return
|
||||
}
|
||||
beg := (n - 1) * 12
|
||||
ckeyLast = binary.LittleEndian.Uint64(r.headers[beg : beg+8])
|
||||
return
|
||||
}
|
||||
|
||||
// Len() indicates the total number of containers the iterator expects to have.
|
||||
func (r *baseRoaringIterator) Len() int64 {
|
||||
return r.keys
|
||||
|
|
|
|||
|
|
@ -4430,6 +4430,17 @@ func TestCloneRoaringIterator(t *testing.T) {
|
|||
|
||||
itr2 := itr.Clone()
|
||||
|
||||
firstCkey, lastCkey, empty := itr2.ContainerKeySpan()
|
||||
if empty {
|
||||
t.Fatalf("should not be empty")
|
||||
}
|
||||
if firstCkey != 0 {
|
||||
t.Fatalf("firstCkey should be 0")
|
||||
}
|
||||
if lastCkey != 10001 {
|
||||
t.Fatalf("lastCkey should be 10001")
|
||||
}
|
||||
|
||||
var keys []uint64
|
||||
for itrKey, synthC := itr.NextContainer(); synthC != nil; itrKey, synthC = itr.NextContainer() {
|
||||
keys = append(keys, itrKey)
|
||||
|
|
@ -4446,6 +4457,54 @@ func TestCloneRoaringIterator(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestRoaringIteratorContainerKeySpan(t *testing.T) {
|
||||
|
||||
ca := NewContainerArray([]uint16{1, 10, 100, 1000})
|
||||
ba := NewFileBitmap()
|
||||
ba.Containers.Put(101, ca)
|
||||
ba.Containers.Put(10, ca)
|
||||
ba.Containers.Put(10001, ca)
|
||||
var buf bytes.Buffer
|
||||
_, err := ba.WriteTo(&buf)
|
||||
if err != nil {
|
||||
t.Fatalf("error writing: %v", err)
|
||||
}
|
||||
|
||||
itr, err := NewRoaringIterator(buf.Bytes())
|
||||
if err != nil {
|
||||
t.Fatalf("error NewRoaringIterator(buf.Bytes()): %v", err)
|
||||
}
|
||||
|
||||
firstCkey, lastCkey, empty := itr.ContainerKeySpan()
|
||||
if empty {
|
||||
t.Fatalf("should not be empty")
|
||||
}
|
||||
if firstCkey != 10 {
|
||||
t.Fatalf("firstCkey should be 10")
|
||||
}
|
||||
if lastCkey != 10001 {
|
||||
t.Fatalf("lastCkey should be 10001")
|
||||
}
|
||||
|
||||
// make and check empty bitmap
|
||||
|
||||
baEmpty := NewFileBitmap()
|
||||
var bufEmpty bytes.Buffer
|
||||
_, err = baEmpty.WriteTo(&bufEmpty)
|
||||
if err != nil {
|
||||
t.Fatalf("error writing: %v", err)
|
||||
}
|
||||
|
||||
itrEmpty, err := NewRoaringIterator(bufEmpty.Bytes())
|
||||
if err != nil {
|
||||
t.Fatalf("error NewRoaringIterator(bufEmpty.Bytes()): %v", err)
|
||||
}
|
||||
_, _, empty = itrEmpty.ContainerKeySpan()
|
||||
if !empty {
|
||||
t.Fatalf("should be empty")
|
||||
}
|
||||
}
|
||||
|
||||
// we were seeing unionInterval16InPlace() returning too
|
||||
// large an run container, which was causing problems when
|
||||
// we write to the transactional backends. Verify that
|
||||
|
|
|
|||
|
|
@ -135,12 +135,14 @@ func (h *GRPCHandler) execSQL(ctx context.Context, queryStr string) (pb.ToRowser
|
|||
|
||||
// QuerySQL handles the SQL request and sends RowResponses to the stream.
|
||||
func (h *GRPCHandler) QuerySQL(req *pb.QuerySQLRequest, stream pb.Pilosa_QuerySQLServer) error {
|
||||
start := time.Now()
|
||||
results, err := h.execSQL(stream.Context(), req.Sql)
|
||||
duration := time.Since(start)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
err = results.ToRows(stream.Send)
|
||||
err = newDurationRowser(results, duration).ToRows(stream.Send)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "streaming result")
|
||||
}
|
||||
|
|
@ -161,14 +163,24 @@ func (h *GRPCHandler) QuerySQL(req *pb.QuerySQLRequest, stream pb.Pilosa_QuerySQ
|
|||
// concurrently. There is additional discussion and historical context here:
|
||||
// https://github.com/molecula/pilosa/pull/644
|
||||
func (h *GRPCHandler) QuerySQLUnary(ctx context.Context, req *pb.QuerySQLRequest) (*pb.TableResponse, error) {
|
||||
start := time.Now()
|
||||
results, err := h.execSQL(ctx, req.Sql)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if results, ok := results.(pb.ToTabler); ok {
|
||||
return results.ToTable()
|
||||
|
||||
var table *pb.TableResponse
|
||||
switch results := results.(type) {
|
||||
case pb.ToTabler:
|
||||
table, err = results.ToTable()
|
||||
default:
|
||||
table, err = pb.RowsToTable(results, 0)
|
||||
}
|
||||
return pb.RowsToTable(results, 0)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
table.Duration = int64(time.Since(start))
|
||||
return table, nil
|
||||
}
|
||||
|
||||
// QueryPQL handles the PQL request and sends RowResponses to the stream.
|
||||
|
|
@ -200,7 +212,7 @@ func (h *GRPCHandler) QueryPQL(req *pb.QueryPQLRequest, stream pb.Pilosa_QueryPQ
|
|||
}
|
||||
|
||||
t = time.Now()
|
||||
if err := toRowser.ToRows(stream.Send); err != nil {
|
||||
if err := newDurationRowser(toRowser, durQuery).ToRows(stream.Send); err != nil {
|
||||
return errToStatusError(err)
|
||||
}
|
||||
durFormat := time.Since(t)
|
||||
|
|
@ -246,6 +258,8 @@ func (h *GRPCHandler) QueryPQLUnary(ctx context.Context, req *pb.QueryPQLRequest
|
|||
}
|
||||
durFormat := time.Since(t)
|
||||
|
||||
table.Duration = int64(durQuery + durFormat)
|
||||
|
||||
h.stats.Timing(pilosa.MetricGRPCUnaryQueryDurationSeconds, durQuery, 0.1)
|
||||
h.stats.Timing(pilosa.MetricGRPCUnaryFormatDurationSeconds, durFormat, 0.1)
|
||||
h.stats.Count(pilosa.MetricPqlQueries, 1, 1)
|
||||
|
|
@ -430,6 +444,31 @@ func ToRowserWrapper(result interface{}) (pb.ToRowser, error) {
|
|||
return toRowser, nil
|
||||
}
|
||||
|
||||
// durationRowser is a wrapper for pb.ToRowser that can be used to inject a
|
||||
// duration value into the first record in a stream
|
||||
type durationRowser struct {
|
||||
pb.ToRowser
|
||||
duration time.Duration
|
||||
once sync.Once
|
||||
}
|
||||
|
||||
func (r *durationRowser) ToRows(callback func(*pb.RowResponse) error) error {
|
||||
cb := func(rr *pb.RowResponse) error {
|
||||
r.once.Do(func() {
|
||||
rr.Duration = int64(r.duration)
|
||||
})
|
||||
return callback(rr)
|
||||
}
|
||||
return r.ToRowser.ToRows(cb)
|
||||
}
|
||||
|
||||
func newDurationRowser(orig pb.ToRowser, duration time.Duration) pb.ToRowser {
|
||||
return &durationRowser{
|
||||
ToRowser: orig,
|
||||
duration: duration,
|
||||
}
|
||||
}
|
||||
|
||||
// Inspect handles the inspect request and sends an InspectResponse to the stream.
|
||||
func (h *GRPCHandler) Inspect(req *pb.InspectRequest, stream pb.Pilosa_InspectServer) error {
|
||||
const defaultLimit = 100000
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@ import (
|
|||
"fmt"
|
||||
"reflect"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/pilosa/pilosa/v2"
|
||||
|
|
@ -347,7 +348,7 @@ func TestQueryPQLUnary(t *testing.T) {
|
|||
ctx := context.Background()
|
||||
gh := server.NewGRPCHandler(m.API)
|
||||
|
||||
_, err := gh.QueryPQLUnary(ctx, &pb.QueryPQLRequest{
|
||||
resp, err := gh.QueryPQLUnary(ctx, &pb.QueryPQLRequest{
|
||||
Index: i.Name(),
|
||||
Pql: `Set(0, f="zero")`,
|
||||
})
|
||||
|
|
@ -356,6 +357,10 @@ func TestQueryPQLUnary(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if resp.Duration == 0 {
|
||||
t.Fatal("duration not recorded")
|
||||
}
|
||||
|
||||
_, err = gh.QueryPQLUnary(ctx, &pb.QueryPQLRequest{
|
||||
Index: i.Name(),
|
||||
Pql: `Set(1, f="one") Set(2, f="two")`,
|
||||
|
|
@ -367,6 +372,89 @@ func TestQueryPQLUnary(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestQueryPQL(t *testing.T) {
|
||||
// TODO: Replace TestQueryPQL and TestQueryPQLUnary with table-driven test, like TestQuerySQL
|
||||
m := test.RunCommand(t)
|
||||
defer m.Close()
|
||||
|
||||
i := m.MustCreateIndex(t, "i", pilosa.IndexOptions{Keys: false, TrackExistence: true})
|
||||
m.MustCreateField(t, i.Name(), "f", pilosa.OptFieldKeys())
|
||||
gh := server.NewGRPCHandler(m.API)
|
||||
|
||||
mock := &mockPilosa_QuerySQLServer{}
|
||||
|
||||
err := gh.QueryPQL(&pb.QueryPQLRequest{
|
||||
Index: i.Name(),
|
||||
Pql: `Set(0, f="zero")`,
|
||||
}, mock)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if len(mock.Results) != 1 {
|
||||
t.Fatal("expecting one result")
|
||||
}
|
||||
|
||||
if len(mock.Results[0].Headers) != 1 {
|
||||
t.Fatal("expecting one header")
|
||||
}
|
||||
|
||||
if len(mock.Results[0].Columns) != 1 {
|
||||
t.Fatal("expecting one column")
|
||||
}
|
||||
|
||||
if mock.Results[0].Duration == 0 {
|
||||
t.Fatal("expecting non-zero duration")
|
||||
}
|
||||
|
||||
// Set second value so that All() returns more than one result
|
||||
err = gh.QueryPQL(&pb.QueryPQLRequest{
|
||||
Index: i.Name(),
|
||||
Pql: `Set(1, f="zero")`,
|
||||
}, mock)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
mock.clearResults()
|
||||
|
||||
err = gh.QueryPQL(&pb.QueryPQLRequest{
|
||||
Index: i.Name(),
|
||||
Pql: `All()`,
|
||||
}, mock)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if len(mock.Results) != 2 {
|
||||
t.Fatal("expecting two results")
|
||||
}
|
||||
|
||||
if len(mock.Results[0].Headers) != 1 {
|
||||
t.Fatal("expecting one header")
|
||||
}
|
||||
|
||||
if len(mock.Results[0].Columns) != 1 {
|
||||
t.Fatal("expecting one column")
|
||||
}
|
||||
|
||||
if len(mock.Results[1].Headers) != 0 {
|
||||
t.Fatal("expecting no headers on second result")
|
||||
}
|
||||
|
||||
if len(mock.Results[1].Columns) != 1 {
|
||||
t.Fatal("expecting one column on second result")
|
||||
}
|
||||
|
||||
if mock.Results[0].Duration == 0 {
|
||||
t.Fatal("expecting non-zero duration")
|
||||
}
|
||||
|
||||
if mock.Results[1].Duration != 0 {
|
||||
t.Fatal("expecting zero duration on second result")
|
||||
}
|
||||
}
|
||||
|
||||
type (
|
||||
tableResponse struct {
|
||||
headers []columnInfo
|
||||
|
|
@ -382,7 +470,7 @@ type (
|
|||
columnResponse interface{}
|
||||
)
|
||||
|
||||
func TestQuerySQLUnary(t *testing.T) {
|
||||
func TestQuerySQL(t *testing.T) {
|
||||
|
||||
ctx := context.Background()
|
||||
gh, tearDownFunc := setUpTestQuerySQLUnary(ctx, t)
|
||||
|
|
@ -824,12 +912,33 @@ func TestQuerySQLUnary(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("sql: %s, error: %v", test.sql, err)
|
||||
} else {
|
||||
if resp.Duration == 0 {
|
||||
t.Fatal("duration not recorded")
|
||||
}
|
||||
tr := toTableResponse(resp)
|
||||
if err := test.eq(test.exp, tr); err != nil {
|
||||
t.Fatalf("sql: %s, error: %+v", test.sql, err)
|
||||
}
|
||||
}
|
||||
})
|
||||
t.Run("test-"+strconv.Itoa(i)+"-streaming", func(t *testing.T) {
|
||||
if strings.HasPrefix(test.sql, "drop table") {
|
||||
t.Skip("drop statements can only run once")
|
||||
}
|
||||
mock := &mockPilosa_QuerySQLServer{}
|
||||
err := gh.QuerySQL(&pb.QuerySQLRequest{Sql: test.sql}, mock)
|
||||
if err != nil {
|
||||
t.Fatalf("sql: %s, error: %v", test.sql, err)
|
||||
} else {
|
||||
if mock.Results[0].Duration == 0 {
|
||||
t.Fatal("duration not recorded")
|
||||
}
|
||||
if len(mock.Results) > 1 && mock.Results[1].Duration != 0 {
|
||||
t.Fatal("duration on second result expected to be zero")
|
||||
}
|
||||
// TODO: test result values
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1099,3 +1208,21 @@ func equalUnordered(exp tableResponse, got tableResponse) error {
|
|||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
type mockPilosa_QuerySQLServer struct {
|
||||
pb.Pilosa_QuerySQLServer
|
||||
Results []*pb.RowResponse
|
||||
}
|
||||
|
||||
func (m *mockPilosa_QuerySQLServer) Send(result *pb.RowResponse) error {
|
||||
m.Results = append(m.Results, result)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockPilosa_QuerySQLServer) Context() context.Context {
|
||||
return context.Background()
|
||||
}
|
||||
|
||||
func (m *mockPilosa_QuerySQLServer) clearResults() {
|
||||
m.Results = m.Results[:0]
|
||||
}
|
||||
|
|
|
|||
32
txfactory.go
32
txfactory.go
|
|
@ -38,7 +38,6 @@ import (
|
|||
// public strings that pilosa/server/config.go can reference
|
||||
const (
|
||||
RoaringTxn string = "roaring"
|
||||
LmdbTxn string = "lmdb"
|
||||
RBFTxn string = "rbf"
|
||||
BoltTxn string = "bolt"
|
||||
)
|
||||
|
|
@ -53,7 +52,7 @@ const DefaultTxsrc = RBFTxn
|
|||
// which the transaction has committed or rolled back. Since
|
||||
// memory segments will be recycled by the underlying databases,
|
||||
// this can lead to corruption. When DetectMemAccessPastTx is true,
|
||||
// code in lmdb.go will copy the transactionally viewed memory before
|
||||
// code in bolt.go will copy the transactionally viewed memory before
|
||||
// returning it for bitmap reading, and then zero it or overwrite it
|
||||
// with -2 when the Tx completes.
|
||||
//
|
||||
|
|
@ -83,8 +82,7 @@ var sep = string(os.PathSeparator)
|
|||
// course that your "new" read Tx actually has an "old" view
|
||||
// of the database.
|
||||
//
|
||||
// At the moment, given that LMDB demands that
|
||||
// all write Tx are created and executed on the same C thread, most
|
||||
// At the moment, most
|
||||
// writes to individual shards are commited eagerly and locally
|
||||
// when the `defer finisher(&err0)` is run.
|
||||
// This is done by returning a finisher that actually Commits,
|
||||
|
|
@ -265,16 +263,8 @@ func (qcx *Qcx) GetTx(o Txo) (tx Tx, finisher func(perr *error), err error) {
|
|||
// don't deadlock against themselves under blue-green.
|
||||
o.Write = o.Write || qcx.write
|
||||
|
||||
// note: write Tx were re-using Tx across different goroutines,
|
||||
// which lmdb will not be pleased with. For reads this
|
||||
// should be okay, as the docs say
|
||||
// "If you want to pass read-only transactions across threads,
|
||||
// you can use the MDB_NOTLS option on the environment."
|
||||
// -- http://www.lmdb.tech/doc/starting.html
|
||||
// and we always use lmdb.NoTLS as the lmdb-go bindings ensure this.
|
||||
//
|
||||
// So we make ALL write transactions local, and never reuse them
|
||||
// below.
|
||||
// In general, we make ALL write transactions local, and never reuse them
|
||||
// below. Previously this was to help lmdb.
|
||||
//
|
||||
// *However* there is one exception: when we have set RequiredForAtomicWriteTx
|
||||
// for the importing of an AtomicRequest, then we must use that
|
||||
|
|
@ -400,7 +390,7 @@ func (qcx *Qcx) ListOpenTx() string {
|
|||
}
|
||||
|
||||
// TxFactory abstracts the creation of Tx interface-level
|
||||
// transactions so that RBF, BoltDB, LMDB, or Roaring-fragment-files, or several
|
||||
// transactions so that RBF, BoltDB, or Roaring-fragment-files, or several
|
||||
// of these at once in parallel, is used as the storage and transction layer.
|
||||
type TxFactory struct {
|
||||
typeOfTx string
|
||||
|
|
@ -435,13 +425,12 @@ const (
|
|||
noneTxn txtype = 0
|
||||
roaringTxn txtype = 1 // these don't really have any transactions
|
||||
rbfTxn txtype = 2
|
||||
lmdbTxn txtype = 3
|
||||
boltTxn txtype = 4
|
||||
)
|
||||
|
||||
// these need to be skipped by the holder.go field scanner that
|
||||
// calls IsTxDatabasePath
|
||||
var allTypesWithSuffixes = []txtype{rbfTxn, lmdbTxn, boltTxn}
|
||||
var allTypesWithSuffixes = []txtype{rbfTxn, boltTxn}
|
||||
|
||||
// FileSuffix is used to determine backend directory names.
|
||||
// We append '@' to be sure we never collide with a field name
|
||||
|
|
@ -454,8 +443,6 @@ func (ty txtype) FileSuffix() string {
|
|||
return ""
|
||||
case rbfTxn:
|
||||
return "-rbfdb@"
|
||||
case lmdbTxn:
|
||||
return "-lmdb@"
|
||||
case boltTxn:
|
||||
return "-boltdb@"
|
||||
}
|
||||
|
|
@ -504,8 +491,6 @@ func MustTxsrcToTxtype(txsrc string) (types []txtype) {
|
|||
types = append(types, roaringTxn)
|
||||
case RBFTxn: // "rbf"
|
||||
types = append(types, rbfTxn)
|
||||
case LmdbTxn: // "lmdb"
|
||||
types = append(types, lmdbTxn)
|
||||
case BoltTxn: // "bolt"
|
||||
types = append(types, boltTxn)
|
||||
default:
|
||||
|
|
@ -886,8 +871,6 @@ func (ty txtype) String() string {
|
|||
return "roaring"
|
||||
case rbfTxn:
|
||||
return "rbf"
|
||||
case lmdbTxn:
|
||||
return "lmdb"
|
||||
case boltTxn:
|
||||
return "bolt"
|
||||
}
|
||||
|
|
@ -1253,9 +1236,6 @@ func anyGlobalDBWrappersStillOpen() bool {
|
|||
if globalRbfDBReg.Size() != 0 {
|
||||
return true
|
||||
}
|
||||
if globalLMDBReg.Size() != 0 {
|
||||
return true
|
||||
}
|
||||
if globalBoltReg.Size() != 0 {
|
||||
return true
|
||||
}
|
||||
|
|
|
|||
|
|
@ -22,7 +22,6 @@ import (
|
|||
"time"
|
||||
|
||||
"github.com/glycerine/lmdb-go/lmdb"
|
||||
//"github.com/pilosa/pilosa/v2/logger"
|
||||
)
|
||||
|
||||
func Test_TxFactory_Qcx_query_context(t *testing.T) {
|
||||
|
|
@ -121,7 +120,7 @@ func Test_TxFactory_UpdateBlueFromGreen_OnStartup(t *testing.T) {
|
|||
orig := os.Getenv("PILOSA_TXSRC")
|
||||
defer os.Setenv("PILOSA_TXSRC", orig) // must restore or will mess up other tests!
|
||||
|
||||
checked := []string{"lmdb", "roaring", "rbf"}
|
||||
checked := []string{"roaring", "rbf"}
|
||||
|
||||
expectError := false
|
||||
for _, blue := range checked {
|
||||
|
|
@ -270,7 +269,7 @@ func Test_TxFactory_verifyBlueEqualsGreen(t *testing.T) {
|
|||
orig := os.Getenv("PILOSA_TXSRC")
|
||||
defer os.Setenv("PILOSA_TXSRC", orig) // must restore or will mess up other tests!
|
||||
|
||||
checked := []string{"lmdb", "roaring", "bolt", "rbf"}
|
||||
checked := []string{"roaring", "bolt", "rbf"}
|
||||
|
||||
for _, blue := range checked {
|
||||
for _, green := range checked {
|
||||
|
|
@ -408,8 +407,8 @@ func Test_TxFactory_verifyStringConstantsMatch(t *testing.T) {
|
|||
// our const definitions at the top of txfactory.go, or
|
||||
// else blue-green transactions cannot determine when
|
||||
// the second transaction is being released in dbshard.go.
|
||||
check := []txtype{roaringTxn, rbfTxn, lmdbTxn, boltTxn}
|
||||
expect := []string{RoaringTxn, RBFTxn, LmdbTxn, BoltTxn}
|
||||
check := []txtype{roaringTxn, rbfTxn, boltTxn}
|
||||
expect := []string{RoaringTxn, RBFTxn, BoltTxn}
|
||||
for i, chk := range check {
|
||||
obs := chk.String()
|
||||
if obs != expect[i] {
|
||||
|
|
|
|||
|
|
@ -51,6 +51,7 @@ func init() {
|
|||
// keeper linter happy
|
||||
_ = pp
|
||||
_ = vv
|
||||
_ = DirExists
|
||||
}
|
||||
|
||||
func PanicOn(err error) {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue