mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 02:44:59 +00:00
* adding more error handling for rbf
Co-authored-by: Jacob Brinlee <jacobbrinlee@Jacobs-MBP.attlocal.net>
(cherry picked from commit c9ce26ce96)
2433 lines
60 KiB
Go
2433 lines
60 KiB
Go
// Copyright 2022 Molecula Corp. (DBA FeatureBase).
|
|
// SPDX-License-Identifier: Apache-2.0
|
|
package rbf
|
|
|
|
import (
|
|
"bufio"
|
|
"fmt"
|
|
"io"
|
|
"sort"
|
|
"strings"
|
|
"sync"
|
|
|
|
"github.com/benbjohnson/immutable"
|
|
"github.com/featurebasedb/featurebase/v3/roaring"
|
|
txkey "github.com/featurebasedb/featurebase/v3/short_txkey"
|
|
"github.com/featurebasedb/featurebase/v3/vprint"
|
|
"github.com/pkg/errors"
|
|
)
|
|
|
|
var _ = txkey.ToString
|
|
|
|
// Tx represents an RBF transaction. Transactions provide guarantees such as
|
|
// atomicity for all writes that occur as well as serializable isolation.
|
|
// Transactions can be obtained by calling DB.Begin() and provide a snapshot
|
|
// view at the point-in-time they are started.
|
|
type Tx struct {
|
|
mu sync.RWMutex
|
|
db *DB // parent db
|
|
meta [PageSize]byte // copy of current meta page
|
|
walID int64 // max WAL ID at start of tx
|
|
walPageN int // wal page count
|
|
rootRecords *immutable.SortedMap[string, uint32] // read-only cache of root records
|
|
|
|
// pageMap holds WAL pages that have not yet been transferred
|
|
// into the database pages. So it can be empty, if the whole previous
|
|
// WAL has been checkpointed back into the database.
|
|
pageMap *PageMap // mapping of database pages to WAL IDs
|
|
writable bool // if true, tx can write
|
|
|
|
dirtyPages map[uint32][]byte // updated pages in this tx
|
|
dirtyBitmapPages map[uint32][]byte // updated bitmap pages in this tx
|
|
|
|
// If Rollback() has already completed, don't do it again.
|
|
// Note db == nil means that commit has already been done.
|
|
rollbackDone bool
|
|
|
|
// DeleteEmptyContainer lets us by default match the roaring
|
|
// behavior where an existing container has all its bits cleared
|
|
// but still sticks around in the database.
|
|
DeleteEmptyContainer bool
|
|
|
|
// It is possible for a modification of the free list to cause a page to
|
|
// be allocated or deallocated, which would modify the free list.
|
|
//
|
|
// For the case where pages need to be allocated during free list
|
|
// changes, we can trivially just allocate new pages and not use the
|
|
// free list. That's the simple case...
|
|
modifyingFreelist bool
|
|
// But removals can't be deferred/not-done like that. If a free list
|
|
// change causes us to deallocate a page (such as if we're *allocating*
|
|
// a page, which causes it to be *removed* from the free list), we really
|
|
// do need to record that, but if we try to do it during the update
|
|
// process, things could go horribly wrong. So, we have a transient list
|
|
// of page numbers which have been deallocated, but it happened *during*
|
|
// the modification of the free list. The top-level modification then
|
|
// processes them on its way out, using a defer. During the processing
|
|
// of this list, we *still* have the flag set, and we make a new list
|
|
// while processing the list, so if somehow a pending add to the list
|
|
// manages to trigger a *deallocation* (which I don't think should be
|
|
// happening), we'll process that one after the current list is processed.
|
|
pendingFreelistAdds []uint32
|
|
|
|
// DEBUG
|
|
stack []byte
|
|
}
|
|
|
|
// DBPath returns the path to the directory that holds the parent database.
|
|
func (tx *Tx) DBPath() string {
|
|
return tx.db.Path
|
|
}
|
|
|
|
// Writable returns true if the transaction can mutate data. Using transaction
|
|
// methods that attempt to write will return ErrTxNotWritable.
|
|
func (tx *Tx) Writable() bool {
|
|
return tx.writable
|
|
}
|
|
|
|
// dirty returns true if any pages have been updated in this tx.
|
|
func (tx *Tx) dirty() bool {
|
|
return tx.dirtyN() != 0
|
|
}
|
|
|
|
// dirtyN returns the number of dirty pages.
|
|
func (tx *Tx) dirtyN() int {
|
|
return len(tx.dirtyPages) + len(tx.dirtyBitmapPages)
|
|
}
|
|
|
|
// PageN returns the number of pages in the database as seen by this transaction.
|
|
func (tx *Tx) PageN() int {
|
|
return int(readMetaPageN(tx.meta[:]))
|
|
}
|
|
|
|
// Commit completes the transaction and persists data changes. If this method
|
|
// fails, changes may or may not have been persisted to disk. If no changes have
|
|
// been made during the transaction, this functions the same as a rollback.
|
|
//
|
|
// Attempting to commit an already committed or rolled back transaction will
|
|
// return an ErrTxClosed error.
|
|
func (tx *Tx) Commit() error {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
if tx.db == nil {
|
|
return ErrTxClosed
|
|
}
|
|
|
|
// Remove any free pages off the end of the file and update the size.
|
|
if err := tx.truncateFreelist(); err != nil {
|
|
return err
|
|
}
|
|
|
|
// If any pages have been written, ensure we write a new meta page with
|
|
// the commit flag to mark the end of the transaction.
|
|
if tx.dirty() {
|
|
if err := tx.flush(); err != nil {
|
|
return err
|
|
}
|
|
// future plan: after checkpoint is moved to background
|
|
// or not every removeTx, then we can move the
|
|
// tx.db.rootRecords = tx.rootRecords into removeTx().
|
|
//
|
|
// ... or maybe not: let's do that part here, and then removeTx
|
|
// may or may not start a checkpoint, possibly asynchronously.
|
|
//
|
|
// avoid race detector firing on a write race here
|
|
// vs the read of rootRecords at db.Begin(), then release
|
|
// the lock, because we need removeTx to grab the lock to
|
|
// work, but if it wants to checkpoint, it wants to be able to return
|
|
// to us here and still be holding the lock.
|
|
tx.db.mu.Lock()
|
|
tx.db.rootRecords = tx.rootRecords
|
|
tx.db.pageMap = tx.pageMap
|
|
tx.db.walPageN = tx.walPageN
|
|
tx.db.mu.Unlock()
|
|
}
|
|
|
|
tx.db.mu.Lock()
|
|
defer tx.db.mu.Unlock()
|
|
// Disconnect transaction from DB.
|
|
return tx.db.removeTx(tx)
|
|
}
|
|
|
|
// truncateFreelist removes any free pages off the end of the file and updates
|
|
// the size of the database. This allows the data file to be resized on checkpoint.
|
|
func (tx *Tx) truncateFreelist() error {
|
|
for {
|
|
if truncated, err := tx.truncateLastFreePage(); err != nil {
|
|
return err
|
|
} else if !truncated {
|
|
return nil // no more free pages at end of file, exit
|
|
}
|
|
}
|
|
}
|
|
|
|
// truncateLastFreePage removes the last page from the file if it is a free page.
|
|
// The page count is then decremented to move the high water mark to remove the page.
|
|
// Returns true if a page was removed, otherwise returns false.
|
|
func (tx *Tx) truncateLastFreePage() (truncated bool, outErr error) {
|
|
tx.modifyingFreelist = true
|
|
defer tx.freelistCleanup(&outErr)
|
|
|
|
c := tx.db.getFreelistCursor(tx)
|
|
defer c.unpooledClose()
|
|
if err := c.Last(); err == io.EOF {
|
|
return false, nil
|
|
} else if err != nil {
|
|
return false, err
|
|
}
|
|
|
|
elem := &c.stack.elems[c.stack.top]
|
|
leafPage, _, err := c.tx.readPage(elem.pgno)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
cell := readLeafCell(leafPage, elem.index)
|
|
|
|
// If page number is not the last page then exit.
|
|
pgno := uint32((cell.Key << 16) | uint64(cell.lastValue(tx)))
|
|
pageN := readMetaPageN(tx.meta[:])
|
|
if pgno < pageN-1 {
|
|
return false, nil
|
|
}
|
|
|
|
// Otherwise remove it from the freelist.
|
|
if changed, err := c.Remove(uint64(pgno)); err != nil {
|
|
return false, err
|
|
} else if !changed {
|
|
vprint.PanicOn(fmt.Sprintf("tx.Tx.truncateLastFreePage(): double alloc: %d", pgno))
|
|
}
|
|
|
|
// Decrement the page count in the database.
|
|
writeMetaPageN(tx.meta[:], pageN-1)
|
|
|
|
return true, nil
|
|
}
|
|
|
|
// Rollback discards any changes that have been made by the transaction.
|
|
// A commit or rollback must always be called after a transaction finishes.
|
|
// If this is a writable transaction, the write lock will be released on the DB.
|
|
func (tx *Tx) Rollback() { tx.rollback(false) }
|
|
|
|
func (tx *Tx) rollback(hasDBLock bool) {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
// allow Rollback to be called more than once.
|
|
if tx.rollbackDone {
|
|
return
|
|
}
|
|
tx.rollbackDone = true
|
|
if tx.db == nil {
|
|
// Commit already done.
|
|
return
|
|
}
|
|
|
|
// Disconnect transaction from DB.
|
|
if !hasDBLock {
|
|
tx.db.mu.Lock()
|
|
defer tx.db.mu.Unlock()
|
|
}
|
|
vprint.PanicOn(tx.db.removeTx(tx))
|
|
}
|
|
|
|
// Root returns the root page number for a bitmap. Returns 0 if the bitmap does not exist.
|
|
func (tx *Tx) Root(name string) (uint32, error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
return tx.root(name)
|
|
}
|
|
|
|
func (tx *Tx) root(name string) (uint32, error) {
|
|
// Fetch list of roots (or retrieve from cache).
|
|
records, err := tx.RootRecords()
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
// Lookup record by bitmap name and return its root page.
|
|
pgno, ok := records.Get(name)
|
|
if !ok {
|
|
return 0, ErrBitmapNotFound
|
|
}
|
|
return pgno, nil
|
|
}
|
|
|
|
// BitmapNames returns a list of all bitmap names.
|
|
func (tx *Tx) BitmapNames() ([]string, error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
if tx.db == nil {
|
|
return nil, ErrTxClosed
|
|
}
|
|
|
|
// Read list of root records.
|
|
records, err := tx.RootRecords()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
a := make([]string, 0, records.Len())
|
|
for itr := records.Iterator(); !itr.Done(); {
|
|
k, _, _ := itr.Next()
|
|
a = append(a, k)
|
|
}
|
|
return a, nil
|
|
}
|
|
|
|
// BitmapExist returns true if bitmap exists. Returns an error if name is empty.
|
|
func (tx *Tx) BitmapExists(name string) (bool, error) {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
return tx.bitmapExists(name)
|
|
}
|
|
|
|
func (tx *Tx) bitmapExists(name string) (bool, error) {
|
|
if tx.db == nil {
|
|
return false, ErrTxClosed
|
|
} else if name == "" {
|
|
return false, ErrBitmapNameRequired
|
|
}
|
|
|
|
// Read root records and find entry for bitmap.
|
|
records, err := tx.RootRecords()
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
_, ok := records.Get(name)
|
|
return ok, nil
|
|
}
|
|
|
|
// CreateBitmap creates a new empty bitmap with the given name.
|
|
// Returns an error if the bitmap already exists.
|
|
func (tx *Tx) CreateBitmap(name string) error {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
return tx.createBitmap(name)
|
|
}
|
|
|
|
func (tx *Tx) createBitmap(name string) error {
|
|
if tx.db == nil {
|
|
return ErrTxClosed
|
|
} else if !tx.writable {
|
|
return ErrTxNotWritable
|
|
} else if name == "" {
|
|
return ErrBitmapNameRequired
|
|
}
|
|
|
|
// Read list of root records.
|
|
records, err := tx.RootRecords()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Find btree by name. Exit if already exists.
|
|
if _, ok := records.Get(name); ok {
|
|
return ErrBitmapExists
|
|
}
|
|
|
|
// Allocate new root page.
|
|
pgno, err := tx.allocatePgno()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Write root page.
|
|
page := allocPage()
|
|
writePageNo(page, pgno)
|
|
writeFlags(page, PageTypeLeaf)
|
|
writeCellN(page, 0)
|
|
if err := tx.writePage(page); err != nil {
|
|
return err
|
|
}
|
|
|
|
// Insert into correct index.
|
|
records = records.Set(name, pgno)
|
|
if err := tx.writeRootRecordPages(records); err != nil {
|
|
return fmt.Errorf("write bitmaps: %w", err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// CreateBitmapIfNotExists creates a new empty bitmap with the given name.
|
|
// This is a no-op if the bitmap already exists.
|
|
func (tx *Tx) CreateBitmapIfNotExists(name string) error {
|
|
if err := tx.CreateBitmap(name); err != nil && err != ErrBitmapExists {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (tx *Tx) createBitmapIfNotExists(name string) error {
|
|
if err := tx.createBitmap(name); err != nil && err != ErrBitmapExists {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// DeleteBitmap removes a bitmap with the given name.
|
|
// Returns an error if the bitmap does not exist.
|
|
func (tx *Tx) DeleteBitmap(name string) error {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
if tx.db == nil {
|
|
return ErrTxClosed
|
|
} else if !tx.writable {
|
|
return ErrTxNotWritable
|
|
} else if name == "" {
|
|
return ErrBitmapNameRequired
|
|
}
|
|
|
|
// Read list of root records.
|
|
records, err := tx.RootRecords()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Find btree by name. Exit if it doesn't exist.
|
|
pgno, ok := records.Get(name)
|
|
if !ok {
|
|
return fmt.Errorf("bitmap does not exist: %q", name)
|
|
}
|
|
|
|
// Deallocate all pages in the tree.
|
|
if err := tx.deallocateTree(pgno); err != nil {
|
|
return err
|
|
}
|
|
|
|
// Delete from record list & rewrite record pages.
|
|
records = records.Delete(name)
|
|
if err := tx.writeRootRecordPages(records); err != nil {
|
|
return fmt.Errorf("write bitmaps: %w", err)
|
|
}
|
|
tx.rootRecords = records
|
|
return nil
|
|
}
|
|
|
|
// DeleteBitmapsWithPrefix removes all bitmaps with a given prefix.
|
|
func (tx *Tx) DeleteBitmapsWithPrefix(prefix string) error {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
if tx.db == nil {
|
|
return ErrTxClosed
|
|
} else if !tx.writable {
|
|
return ErrTxNotWritable
|
|
}
|
|
|
|
// Read list of root records.
|
|
records, err := tx.RootRecords()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
for itr := records.Iterator(); !itr.Done(); {
|
|
name, pgno, _ := itr.Next()
|
|
|
|
// Skip bitmaps without matching prefix.
|
|
if !strings.HasPrefix(name, prefix) {
|
|
continue
|
|
}
|
|
// Deallocate all pages in the tree.
|
|
if err := tx.deallocateTree(pgno); err != nil {
|
|
return err
|
|
}
|
|
|
|
records = records.Delete(name)
|
|
}
|
|
|
|
// Rewrite record pages.
|
|
if err := tx.writeRootRecordPages(records); err != nil {
|
|
return fmt.Errorf("write bitmaps: %w", err)
|
|
}
|
|
tx.rootRecords = records
|
|
return nil
|
|
}
|
|
|
|
// RenameBitmap updates the name of an existing bitmap.
|
|
// Returns an error if the bitmap does not exist.
|
|
func (tx *Tx) RenameBitmap(oldname, newname string) error {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
if tx.db == nil {
|
|
return ErrTxClosed
|
|
} else if !tx.writable {
|
|
return ErrTxNotWritable
|
|
} else if oldname == "" || newname == "" {
|
|
return ErrBitmapNameRequired
|
|
}
|
|
|
|
// Read list of root records.
|
|
records, err := tx.RootRecords()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Find btree by name. Exit if it doesn't exist.
|
|
pgno, ok := records.Get(oldname)
|
|
if !ok {
|
|
return fmt.Errorf("bitmap does not exist: %q", oldname)
|
|
}
|
|
|
|
// Update record name & rewrite record pages.
|
|
records = records.Delete(oldname)
|
|
records = records.Set(newname, pgno)
|
|
if err := tx.writeRootRecordPages(records); err != nil {
|
|
return fmt.Errorf("write bitmaps: %w", err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// RootRecords returns a list of root records.
|
|
func (tx *Tx) RootRecords() (records *immutable.SortedMap[string, uint32], err error) {
|
|
if tx.rootRecords != nil {
|
|
return tx.rootRecords, nil
|
|
}
|
|
|
|
records = immutable.NewSortedMap[string, uint32](nil)
|
|
for pgno := readMetaRootRecordPageNo(tx.meta[:]); pgno != 0; {
|
|
page, _, err := tx.readPage(pgno)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// Read all records on the page.
|
|
a, err := readRootRecords(page)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
for _, rec := range a {
|
|
records = records.Set(rec.Name, rec.Pgno)
|
|
}
|
|
|
|
// Read next overflow page number.
|
|
pgno = WalkRootRecordPages(page)
|
|
}
|
|
|
|
// Cache result
|
|
tx.rootRecords = records
|
|
return records, nil
|
|
}
|
|
|
|
// writeRootRecordPages writes a list of root record pages.
|
|
func (tx *Tx) writeRootRecordPages(records *immutable.SortedMap[string, uint32]) (err error) {
|
|
// Release all existing root record pages.
|
|
for pgno := readMetaRootRecordPageNo(tx.meta[:]); pgno != 0; {
|
|
page, _, err := tx.readPage(pgno)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
err = tx.freePgno(pgno)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
pgno = WalkRootRecordPages(page)
|
|
}
|
|
|
|
// Exit early if no records exist.
|
|
if records.Len() == 0 {
|
|
writeMetaRootRecordPageNo(tx.meta[:], 0)
|
|
return nil
|
|
}
|
|
|
|
// Allocate initial root record page.
|
|
pgno, err := tx.allocatePgno()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
writeMetaRootRecordPageNo(tx.meta[:], pgno)
|
|
|
|
// Write new root record pages.
|
|
for itr := records.Iterator(); !itr.Done(); {
|
|
// Initialize page & write as many records as will fit.
|
|
page := allocPage()
|
|
writePageNo(page, pgno)
|
|
writeFlags(page, PageTypeRootRecord)
|
|
|
|
if err := writeRootRecords(page, itr); err == io.ErrShortBuffer {
|
|
// Allocate next pgno and write overflow if we have remaining records.
|
|
if pgno, err = tx.allocatePgno(); err != nil {
|
|
return err
|
|
}
|
|
writeRootRecordOverflowPgno(page, pgno)
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Write page to disk.
|
|
if err := tx.writePage(page); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
// Update cache records.
|
|
tx.rootRecords = records
|
|
|
|
return nil
|
|
}
|
|
|
|
// Add sets a given bit on the bitmap.
|
|
func (tx *Tx) Add(name string, a ...uint64) (changeCount int, err error) {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
if tx.db == nil {
|
|
return 0, ErrTxClosed
|
|
} else if !tx.writable {
|
|
return 0, ErrTxNotWritable
|
|
} else if name == "" {
|
|
return 0, ErrBitmapNameRequired
|
|
}
|
|
|
|
if err := tx.createBitmapIfNotExists(name); err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
c, err := tx.cursor(name)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
defer c.Close()
|
|
|
|
for _, v := range a {
|
|
if vchanged, err := c.Add(v); err != nil {
|
|
return changeCount, err
|
|
} else if vchanged {
|
|
changeCount++
|
|
}
|
|
}
|
|
return changeCount, nil
|
|
}
|
|
|
|
// Remove unsets a given bit on the bitmap.
|
|
func (tx *Tx) Remove(name string, a ...uint64) (changeCount int, err error) {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
if tx.db == nil {
|
|
return 0, ErrTxClosed
|
|
} else if !tx.writable {
|
|
return 0, ErrTxNotWritable
|
|
} else if name == "" {
|
|
return 0, ErrBitmapNameRequired
|
|
}
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return 0, nil
|
|
} else if err != nil {
|
|
return 0, err
|
|
}
|
|
defer c.Close()
|
|
|
|
for _, v := range a {
|
|
if vchanged, err := c.Remove(v); err != nil {
|
|
return changeCount, err
|
|
} else if vchanged {
|
|
changeCount++
|
|
}
|
|
}
|
|
return changeCount, nil
|
|
}
|
|
|
|
// Contains returns true if the given bit is set on the bitmap.
|
|
func (tx *Tx) Contains(name string, v uint64) (bool, error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
if tx.db == nil {
|
|
return false, ErrTxClosed
|
|
} else if name == "" {
|
|
return false, ErrBitmapNameRequired
|
|
}
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return false, nil
|
|
} else if err != nil {
|
|
return false, err
|
|
}
|
|
defer c.Close()
|
|
|
|
return c.Contains(v)
|
|
}
|
|
|
|
// Depth returns the depth of the b-tree for a bitmap.
|
|
func (tx *Tx) Depth(name string) (int, error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
if tx.db == nil {
|
|
return 0, ErrTxClosed
|
|
} else if name == "" {
|
|
return 0, ErrBitmapNameRequired
|
|
}
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return 0, nil
|
|
} else if err != nil {
|
|
return 0, err
|
|
}
|
|
defer c.Close()
|
|
|
|
if err := c.First(); err != nil {
|
|
return 0, err
|
|
}
|
|
return c.stack.top + 1, nil
|
|
}
|
|
|
|
// Cursor returns an instance of a cursor this bitmap.
|
|
func (tx *Tx) Cursor(name string) (*Cursor, error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
return tx.cursor(name)
|
|
}
|
|
|
|
func (tx *Tx) cursor(name string) (*Cursor, error) {
|
|
if tx.db == nil {
|
|
return nil, ErrTxClosed
|
|
} else if name == "" {
|
|
return nil, ErrBitmapNameRequired
|
|
}
|
|
|
|
root, err := tx.root(name)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
c := tx.db.getCursor(tx)
|
|
c.stack.top = 0
|
|
c.stack.elems[0] = stackElem{pgno: root}
|
|
return c, nil
|
|
}
|
|
|
|
// RoaringBitmap returns a bitmap as a Roaring bitmap.
|
|
func (tx *Tx) RoaringBitmap(name string) (*roaring.Bitmap, error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
if tx.db == nil {
|
|
return nil, ErrTxClosed
|
|
} else if name == "" {
|
|
return nil, ErrBitmapNameRequired
|
|
}
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return roaring.NewSliceBitmap(), nil
|
|
} else if err != nil {
|
|
return nil, err
|
|
}
|
|
defer c.Close()
|
|
|
|
other := roaring.NewSliceBitmap()
|
|
if err := c.First(); err == io.EOF {
|
|
return other, nil
|
|
} else if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
for {
|
|
if err := c.Next(); err == io.EOF {
|
|
return other, nil
|
|
} else if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
elem := &c.stack.elems[c.stack.top]
|
|
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))
|
|
}
|
|
}
|
|
|
|
// Container returns a Roaring container by key.
|
|
func (tx *Tx) Container(name string, key uint64) (*roaring.Container, error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
return tx.container(name, key)
|
|
}
|
|
|
|
func (tx *Tx) container(name string, key uint64) (*roaring.Container, error) {
|
|
if tx.db == nil {
|
|
return nil, ErrTxClosed
|
|
} else if name == "" {
|
|
return nil, ErrBitmapNameRequired
|
|
}
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return nil, nil
|
|
} else if err != nil {
|
|
return nil, err
|
|
}
|
|
defer c.Close()
|
|
|
|
if exact, err := c.Seek(key); err != nil || !exact {
|
|
return nil, err
|
|
}
|
|
|
|
elem := &c.stack.elems[c.stack.top]
|
|
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.
|
|
func (tx *Tx) PutContainer(name string, key uint64, ct *roaring.Container) error {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
return tx.putContainer(name, key, ct)
|
|
}
|
|
|
|
func (tx *Tx) putContainer(name string, key uint64, ct *roaring.Container) error {
|
|
if tx.DeleteEmptyContainer && ct.N() == 0 {
|
|
return tx.removeContainer(name, key)
|
|
}
|
|
|
|
cell := ConvertToLeafArgs(key, ct)
|
|
|
|
if err := tx.createBitmapIfNotExists(name); err != nil {
|
|
return err
|
|
}
|
|
|
|
c, err := tx.cursor(name)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer c.Close()
|
|
|
|
if _, err := c.Seek(cell.Key); err != nil {
|
|
return err
|
|
}
|
|
|
|
return c.putLeafCell(cell)
|
|
}
|
|
|
|
func (tx *Tx) putContainerWithCursor(cur *Cursor, key uint64, ct *roaring.Container) error {
|
|
if tx.DeleteEmptyContainer && ct.N() == 0 {
|
|
if exact, err := cur.Seek(key); err != nil || !exact {
|
|
return err
|
|
}
|
|
return cur.deleteLeafCell(key)
|
|
}
|
|
return cur.putLeafCell(ConvertToLeafArgs(key, ct))
|
|
}
|
|
|
|
// RemoveContainer removes a container from the bitmap by key.
|
|
func (tx *Tx) RemoveContainer(name string, key uint64) error {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
return tx.removeContainer(name, key)
|
|
}
|
|
|
|
func (tx *Tx) removeContainer(name string, key uint64) error {
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return nil
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
defer c.Close()
|
|
|
|
if exact, err := c.Seek(key); err != nil || !exact {
|
|
return err
|
|
}
|
|
return c.deleteLeafCell(key)
|
|
}
|
|
|
|
// Check verifies the integrity of the database.
|
|
func (tx *Tx) Check() error {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
if tx.db == nil {
|
|
return ErrTxClosed
|
|
}
|
|
|
|
var errorList ErrorList
|
|
if err := tx.checkPageAllocations(); err != nil {
|
|
errorList.Append(err)
|
|
}
|
|
return errorList.Err()
|
|
}
|
|
|
|
func (tx *Tx) checkPage(pgno, parent, typ uint32) error {
|
|
switch typ {
|
|
case PageTypeBranch:
|
|
return tx.checkBranchPage(pgno, parent, typ)
|
|
default:
|
|
return nil
|
|
}
|
|
}
|
|
|
|
func (tx *Tx) checkBranchPage(pgno, parent, typ uint32) error {
|
|
page, _, err := tx.readPage(pgno)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if readCellN(page) == 0 {
|
|
return fmt.Errorf("branch page %d is empty", pgno)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// checkPageAllocations ensures that all pages are either in-use or on the freelist.
|
|
func (tx *Tx) checkPageAllocations() error {
|
|
var errorList ErrorList
|
|
freePageSet, err := tx.freePageSet()
|
|
if err != nil {
|
|
errorList.Append(err)
|
|
}
|
|
|
|
inusePageSet, err := tx.inusePageSet()
|
|
if err != nil {
|
|
errorList.Append(err)
|
|
}
|
|
|
|
// Iterate over all pages and ensure they are either in-use or free.
|
|
// They should not be BOTH in-use or free or NEITHER in-use or free.
|
|
pageN := readMetaPageN(tx.meta[:])
|
|
for pgno := uint32(1); pgno < pageN; pgno++ {
|
|
_, isInuse := inusePageSet[pgno]
|
|
_, isFree := freePageSet[pgno]
|
|
|
|
if isInuse && isFree {
|
|
errorList.Append(fmt.Errorf("page in-use & free: pgno=%d", pgno))
|
|
continue
|
|
}
|
|
|
|
if !isInuse && !isFree {
|
|
errorList.Append(fmt.Errorf("page not in-use & not free: pgno=%d", pgno))
|
|
continue
|
|
}
|
|
}
|
|
|
|
return errorList.Err()
|
|
}
|
|
|
|
// freePageSet returns the set of pages in the freelist.
|
|
func (tx *Tx) freePageSet() (map[uint32]struct{}, error) {
|
|
var errorList ErrorList
|
|
|
|
m := make(map[uint32]struct{})
|
|
c := Cursor{tx: tx}
|
|
c.stack.elems[0] = stackElem{pgno: readMetaFreelistPageNo(tx.meta[:])}
|
|
if err := c.First(); err == io.EOF {
|
|
return m, nil
|
|
} else if err != nil {
|
|
return m, err
|
|
}
|
|
|
|
for {
|
|
if err := c.Next(); err == io.EOF {
|
|
return m, errorList.Err()
|
|
} else if err != nil {
|
|
errorList.Append(err)
|
|
return m, errorList.Err()
|
|
}
|
|
|
|
elem := &c.stack.elems[c.stack.top]
|
|
leafPage, _, err := c.tx.readPage(elem.pgno)
|
|
if err != nil {
|
|
errorList.Append(fmt.Errorf("cannot read free page: pgno=%d err=%w", elem.pgno, err))
|
|
continue
|
|
}
|
|
|
|
cell := readLeafCell(leafPage, elem.index)
|
|
for _, v := range cell.Values(tx) {
|
|
pgno := uint32((cell.Key << 16) | uint64(v))
|
|
m[pgno] = struct{}{}
|
|
}
|
|
}
|
|
}
|
|
|
|
// inusePageSet returns the set of pages in use by the root records or b-trees.
|
|
func (tx *Tx) inusePageSet() (map[uint32]struct{}, error) {
|
|
var errorList ErrorList
|
|
m := make(map[uint32]struct{})
|
|
m[0] = struct{}{} // meta page
|
|
|
|
// Traverse root record linked list and mark each page as in-use.
|
|
for pgno := readMetaRootRecordPageNo(tx.meta[:]); pgno != 0; {
|
|
m[pgno] = struct{}{}
|
|
|
|
page, _, err := tx.readPage(pgno)
|
|
if err != nil {
|
|
errorList.Append(err)
|
|
break
|
|
}
|
|
pgno = WalkRootRecordPages(page)
|
|
}
|
|
|
|
// Traverse freelist and mark pages as in-use.
|
|
if err := tx.walkTree(readMetaFreelistPageNo(tx.meta[:]), 0, func(pgno, parent, typ uint32, err error) error {
|
|
if err != nil {
|
|
errorList.Append(err)
|
|
return nil
|
|
}
|
|
|
|
m[pgno] = struct{}{}
|
|
if err := tx.checkPage(pgno, parent, typ); err != nil {
|
|
errorList.Append(err)
|
|
}
|
|
return nil
|
|
}); err != nil {
|
|
return m, err
|
|
}
|
|
|
|
// Traverse every b-tree and mark pages as in-use.
|
|
records, err := tx.RootRecords()
|
|
if err != nil {
|
|
errorList.Append(err)
|
|
} else {
|
|
for itr := records.Iterator(); !itr.Done(); {
|
|
_, pgno, _ := itr.Next()
|
|
|
|
if err := tx.walkTree(pgno, 0, func(pgno, parent, typ uint32, err error) error {
|
|
if err != nil {
|
|
errorList.Append(err)
|
|
}
|
|
|
|
m[pgno] = struct{}{}
|
|
if err := tx.checkPage(pgno, parent, typ); err != nil {
|
|
errorList.Append(err)
|
|
}
|
|
return nil
|
|
}); err != nil {
|
|
return m, err
|
|
}
|
|
}
|
|
}
|
|
|
|
return m, errorList.Err()
|
|
}
|
|
|
|
// GetSizeBytesWithPrefix returns the size of bitmaps with a given key prefix.
|
|
func (tx *Tx) GetSizeBytesWithPrefix(prefix string) (n uint64, err error) {
|
|
records, err := tx.RootRecords()
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
// Loop over each bitmap in the database.
|
|
for itr := records.Iterator(); !itr.Done(); {
|
|
name, pgno, _ := itr.Next()
|
|
|
|
// Skip over any bitmaps that don't have a matching prefix.
|
|
if !strings.HasPrefix(name, prefix) {
|
|
continue
|
|
}
|
|
|
|
// Traverse the bitmap's b-tree and count the bytes for each page.
|
|
if err := tx.walkTree(pgno, 0, func(pgno, parent, typ uint32, err error) error {
|
|
n += PageSize
|
|
return err
|
|
}); err != nil {
|
|
return 0, err
|
|
}
|
|
}
|
|
return n, nil
|
|
}
|
|
|
|
// walkTree recursively iterates over a page and all its children.
|
|
func (tx *Tx) walkTree(pgno, parent uint32, fn func(pgno, parent, typ uint32, err error) error) error {
|
|
// Read page and iterate over children.
|
|
page, _, err := tx.readPage(pgno)
|
|
if err != nil {
|
|
return fn(pgno, parent, 0, fmt.Errorf("cannot read page: pgno=%d parent=%d err=%s", pgno, parent, err))
|
|
}
|
|
|
|
switch typ := readFlags(page); typ {
|
|
case PageTypeBranch:
|
|
if err := fn(pgno, parent, typ, nil); err != nil {
|
|
return err
|
|
}
|
|
|
|
for i, n := 0, readCellN(page); i < n; i++ {
|
|
cell := readBranchCell(page, i)
|
|
if err := tx.walkTree(cell.ChildPgno, pgno, fn); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
|
|
case PageTypeLeaf:
|
|
if err := fn(pgno, parent, typ, nil); err != nil {
|
|
return err
|
|
}
|
|
|
|
// Execute callback only for bitmap pages pointed to by this leaf.
|
|
for i, n := 0, readCellN(page); i < n; i++ {
|
|
if cell := readLeafCell(page, i); cell.Type == ContainerTypeBitmapPtr {
|
|
if err := fn(toPgno(cell.Data), pgno, PageTypeBitmap, nil); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
return nil
|
|
|
|
default:
|
|
return fn(pgno, parent, typ, fmt.Errorf("invalid page type: pgno=%d parent=%d type=%d", pgno, parent, typ))
|
|
}
|
|
}
|
|
|
|
// freelistCleanup handles things which we need to add to the free list,
|
|
// because they became free during the process of modifying the free list.
|
|
// It also marks us as done modifying the free list. Expected usage is that
|
|
// you set modifyingFreelist to true, then defer this.
|
|
//
|
|
// if you are modifying the free list, we can't further change the free
|
|
// list during that modification. for page allocations, we can just skip
|
|
// the free list check. for deallocations, though, we do need to mark them
|
|
// as freed at some point. so, if we're modifying the free list when
|
|
// a new freePgno happens, we stash the new pages in here, then apply
|
|
// them afterwards. so far as i know, this can actually only happen
|
|
// during an allocate, when we're removing entries from the free list, and
|
|
// the add path doesn't ever trigger it. so, when we remove entries from
|
|
// the free list, it's possible that doing so frees up pages that were
|
|
// part of the free list, and we then add them. but we don't have to worry
|
|
// about that removing things from the free list, because the add logic
|
|
// already just uses new pages rather than trying to use the free list
|
|
// when it knows the free list is involved.
|
|
//
|
|
// Because this is expected to be used in a defer, instead of returning an
|
|
// error, it will set the error it got the address of to a new error if it
|
|
// encounters one and there wasn't one already.
|
|
func (tx *Tx) freelistCleanup(outErr *error) {
|
|
defer func() {
|
|
// no matter what, we're done with this after this, but we still
|
|
// want it set *while* we do this so nothing we do will have side
|
|
// effects that collide with what we're doing.
|
|
tx.modifyingFreelist = false
|
|
}()
|
|
if len(tx.pendingFreelistAdds) == 0 {
|
|
return
|
|
}
|
|
c := tx.db.getFreelistCursor(tx)
|
|
defer c.unpooledClose()
|
|
for len(tx.pendingFreelistAdds) > 0 {
|
|
var pass []uint32
|
|
pass, tx.pendingFreelistAdds = tx.pendingFreelistAdds, nil
|
|
for _, pgno := range pass {
|
|
if changed, err := c.Add(uint64(pgno)); err != nil {
|
|
if outErr != nil && *outErr == nil {
|
|
*outErr = err
|
|
}
|
|
return
|
|
} else if !changed {
|
|
vprint.PanicOn(fmt.Sprintf("rbf.Tx.freelistCleanup(): double free: %d", pass))
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// allocatePgno returns a page number for a new available page. This page may be
|
|
// pulled from the free list or, if no free pages are available, it will be
|
|
// created by extending the file size.
|
|
//
|
|
// allocatePgno uses the freelist cursor (a shared db-wide thing), and sets
|
|
// the "modifyingFreelist" flag while it's running. If for some reason a
|
|
// modification to the freelist would require a new allocation or free,
|
|
// allocations always just create a new page, and frees are processed later
|
|
// by a separate call through a deferred tx.freelistCleanup().
|
|
func (tx *Tx) allocatePgno() (_ uint32, outErr error) {
|
|
if tx.modifyingFreelist {
|
|
return tx.allocateNewPgno(), nil
|
|
}
|
|
// this serves as a precaution against double-use of the freelist cursor
|
|
// used database-wide. we don't have actual synchronization here because
|
|
// only one write Tx should exist at once and it's not safe to use its
|
|
// write-capable ops concurrently anyway.
|
|
tx.modifyingFreelist = true
|
|
defer tx.freelistCleanup(&outErr)
|
|
c := tx.db.getFreelistCursor(tx)
|
|
defer c.unpooledClose()
|
|
if err := c.First(); err == io.EOF {
|
|
return tx.allocateNewPgno(), nil
|
|
} else if err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
elem := &c.stack.elems[c.stack.top]
|
|
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))
|
|
|
|
if changed, err := c.Remove(uint64(pgno)); err != nil {
|
|
return 0, err
|
|
} else if !changed {
|
|
vprint.PanicOn(fmt.Sprintf("tx.Tx.allocatePgno(): double alloc: %d", pgno))
|
|
}
|
|
return pgno, nil
|
|
}
|
|
|
|
// allocateNewPgno requests a new page unconditionally, ignoring the free list.
|
|
func (tx *Tx) allocateNewPgno() uint32 {
|
|
// Increment the total page count by one and return the last page.
|
|
pgno := readMetaPageN(tx.meta[:])
|
|
writeMetaPageN(tx.meta[:], pgno+1)
|
|
return pgno
|
|
}
|
|
|
|
// deallocate releases a page number to the freelist.
|
|
func (tx *Tx) freePgno(pgno uint32) (outErr error) {
|
|
delete(tx.dirtyPages, pgno)
|
|
delete(tx.dirtyBitmapPages, pgno)
|
|
|
|
if tx.modifyingFreelist {
|
|
tx.pendingFreelistAdds = append(tx.pendingFreelistAdds, pgno)
|
|
return nil
|
|
}
|
|
c := tx.db.getFreelistCursor(tx)
|
|
defer c.unpooledClose()
|
|
|
|
tx.modifyingFreelist = true
|
|
defer tx.freelistCleanup(&outErr)
|
|
if changed, err := c.Add(uint64(pgno)); err != nil {
|
|
return err
|
|
} else if !changed {
|
|
vprint.PanicOn(fmt.Sprintf("rbf.Tx.freePgno(): double free: %d", pgno))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// deallocateTree recursively all pages in a btree.
|
|
func (tx *Tx) deallocateTree(pgno uint32) error {
|
|
page, _, err := tx.readPage(pgno)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
switch typ := readFlags(page); typ {
|
|
case PageTypeBranch:
|
|
for i, n := 0, readCellN(page); i < n; i++ {
|
|
cell := readBranchCell(page, i)
|
|
if err := tx.deallocateTree(cell.ChildPgno); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return tx.freePgno(pgno)
|
|
|
|
case PageTypeLeaf:
|
|
for i, n := 0, readCellN(page); i < n; i++ {
|
|
if cell := readLeafCell(page, i); cell.Type == ContainerTypeBitmapPtr {
|
|
if err := tx.freePgno(toPgno(cell.Data)); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
return tx.freePgno(pgno)
|
|
default:
|
|
return fmt.Errorf("rbf.Tx.deallocateTree(): invalid page type: pgno=%d type=%d", pgno, typ)
|
|
}
|
|
}
|
|
|
|
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[:], false, nil
|
|
}
|
|
|
|
// Verify page number requested is within current size of database.
|
|
pageN := readMetaPageN(tx.meta[:])
|
|
if pgno >= pageN {
|
|
return nil, false, fmt.Errorf("rbf: page read out of bounds: pgno=%d max=%d", pgno, pageN-1)
|
|
}
|
|
|
|
// Check if page has been updated in this tx.
|
|
if tx.writable {
|
|
if page := tx.dirtyPages[pgno]; page != nil {
|
|
return page, true, nil
|
|
} else if page := tx.dirtyBitmapPages[pgno]; page != nil {
|
|
return page, true, nil
|
|
}
|
|
}
|
|
|
|
// Check if page is remapped in WAL.
|
|
if walID, ok := tx.pageMap.Get(pgno); ok {
|
|
buf, err := tx.db.readWALPageByID(walID)
|
|
return buf, false, err
|
|
}
|
|
|
|
// Otherwise read directly from DB.
|
|
buf, err := tx.db.readDBPage(pgno)
|
|
return buf, false, err
|
|
}
|
|
|
|
func (tx *Tx) writePage(page []byte) error {
|
|
tx.dirtyPages[readPageNo(page)] = page
|
|
return tx.checkTxSize()
|
|
}
|
|
|
|
func (tx *Tx) writeBitmapPage(pgno uint32, page []byte) error {
|
|
tx.dirtyBitmapPages[pgno] = page
|
|
return tx.checkTxSize()
|
|
}
|
|
|
|
func (tx *Tx) checkTxSize() error {
|
|
pageN := tx.walPageN + len(tx.dirtyPages) + (len(tx.dirtyBitmapPages) * 2)
|
|
if pageN*PageSize >= len(tx.db.wal) {
|
|
return ErrTxTooLarge
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (tx *Tx) AddRoaring(name string, bm *roaring.Bitmap) (changed bool, err error) {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
if err := tx.createBitmapIfNotExists(name); err != nil {
|
|
return false, err
|
|
}
|
|
|
|
c, err := tx.cursor(name)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
defer c.Close()
|
|
|
|
return c.AddRoaring(bm)
|
|
}
|
|
|
|
func (tx *Tx) leafCellBitmap(pgno uint32) (uint32, []uint64, error) {
|
|
page, _, err := tx.readPage(pgno)
|
|
if err != nil {
|
|
return 0, nil, err
|
|
}
|
|
return pgno, toArray64(page), err
|
|
}
|
|
|
|
// leafCellBitmapInto copies the bitmap into provided space
|
|
func (tx *Tx) leafCellBitmapInto(pgno uint32, into []byte) (uint32, []uint64, error) {
|
|
page, _, err := tx.readPage(pgno)
|
|
if err != nil {
|
|
return 0, nil, err
|
|
}
|
|
copy(into, page)
|
|
return pgno, toArray64(into), err
|
|
}
|
|
|
|
func (tx *Tx) ContainerIterator(name string, key uint64) (citer roaring.ContainerIterator, found bool, err error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return &emptyContainerIterator{}, false, nil // nothing available.
|
|
} else if err != nil {
|
|
return nil, false, err
|
|
}
|
|
|
|
exact, err := c.Seek(key)
|
|
if err != nil {
|
|
return nil, false, err
|
|
}
|
|
return &containerIterator{cursor: c}, exact, nil
|
|
}
|
|
|
|
// Shared pool for container filters, used because they contain cursors
|
|
// which are large.
|
|
var containerFilterPool = &sync.Pool{}
|
|
|
|
// getContainerFilter generates a containerFilter, which may actually secretly
|
|
// be used for rewriting; the data structures are similar enough that sharing
|
|
// a pool for both types seems advantageous.
|
|
func getContainerFilter(c *Cursor, name string, filter roaring.BitmapFilter, rewriter roaring.BitmapRewriter, tx *Tx) *containerFilter {
|
|
existing := containerFilterPool.Get()
|
|
if existing == nil {
|
|
return &containerFilter{cursor: c, name: name, filter: filter, rewriter: rewriter, tx: tx}
|
|
}
|
|
f := existing.(*containerFilter)
|
|
f.cursor = c
|
|
f.name = name
|
|
f.filter = filter
|
|
f.rewriter = rewriter
|
|
f.tx = tx
|
|
return f
|
|
}
|
|
|
|
func (tx *Tx) ApplyRewriter(name string, key uint64, rewriter roaring.BitmapRewriter) (err error) {
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
// Unlike a Filter, Rewriter makes sense to apply to an empty bitmap.
|
|
if err = tx.createBitmapIfNotExists(name); err != nil {
|
|
return err
|
|
}
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return nil // nothing available.
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
|
|
_, err = c.Seek(key)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
f := getContainerFilter(c, name, nil, rewriter, tx)
|
|
defer f.Close()
|
|
return f.ApplyRewriter()
|
|
}
|
|
|
|
func (tx *Tx) ApplyFilter(name string, key uint64, filter roaring.BitmapFilter) (err error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return nil // nothing available.
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
|
|
_, err = c.Seek(key)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
f := getContainerFilter(c, name, filter, nil, tx)
|
|
defer f.Close()
|
|
return f.ApplyFilter()
|
|
}
|
|
|
|
func (tx *Tx) Count(name string) (uint64, error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return 0, nil
|
|
} else if err != nil {
|
|
return 0, err
|
|
}
|
|
defer c.Close()
|
|
|
|
if err := c.First(); err == io.EOF {
|
|
return 0, nil
|
|
} else if err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
var n uint64
|
|
for {
|
|
if err := c.Next(); err == io.EOF {
|
|
break
|
|
} else if err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
elem := &c.stack.elems[c.stack.top]
|
|
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
|
|
}
|
|
|
|
func (tx *Tx) Max(name string) (uint64, error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return 0, nil
|
|
} else if err != nil {
|
|
return 0, err
|
|
}
|
|
defer c.Close()
|
|
|
|
if err := c.Last(); err == io.EOF {
|
|
return 0, nil
|
|
} else if err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
elem := &c.stack.elems[c.stack.top]
|
|
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
|
|
}
|
|
|
|
func (tx *Tx) Min(name string) (uint64, bool, error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return 0, false, nil
|
|
} else if err != nil {
|
|
return 0, false, err
|
|
}
|
|
defer c.Close()
|
|
|
|
if err := c.First(); err == io.EOF {
|
|
return 0, false, nil
|
|
} else if err != nil {
|
|
return 0, false, err
|
|
}
|
|
|
|
elem := &c.stack.elems[c.stack.top]
|
|
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
|
|
}
|
|
|
|
// roaring.countRange counts the number of bits set between [start, end).
|
|
func (tx *Tx) CountRange(name string, start, end uint64) (uint64, error) {
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
if start >= end {
|
|
return 0, nil
|
|
}
|
|
|
|
skey := highbits(start)
|
|
ekey := highbits(end)
|
|
ebits := int32(lowbits(end))
|
|
|
|
csr, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return 0, nil
|
|
} else if err != nil {
|
|
return 0, err
|
|
}
|
|
defer csr.Close()
|
|
|
|
exact, err := csr.Seek(skey)
|
|
_ = exact
|
|
if err == io.EOF {
|
|
return 0, nil
|
|
} else if err != nil {
|
|
return 0, err
|
|
}
|
|
var n uint64
|
|
for {
|
|
if err := csr.Next(); err == io.EOF {
|
|
break
|
|
} else if err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
elem := &csr.stack.elems[csr.stack.top]
|
|
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
|
|
}
|
|
|
|
// If range is entirely in one container then just count that range.
|
|
if skey == ekey {
|
|
return uint64(c.countRange(tx, int32(lowbits(start)), ebits)), nil
|
|
}
|
|
// INVAR: skey < ekey
|
|
|
|
// k > ekey handles the case when start > end and where start and end
|
|
// are in different containers. Same container case is already handled above.
|
|
if k > ekey {
|
|
break
|
|
}
|
|
if k == skey {
|
|
n += uint64(c.countRange(tx, int32(lowbits(start)), roaring.MaxContainerVal+1))
|
|
continue
|
|
}
|
|
if k < ekey {
|
|
n += uint64(c.BitN)
|
|
continue
|
|
}
|
|
if k == ekey && ebits > 0 {
|
|
n += uint64(c.countRange(tx, 0, ebits))
|
|
break
|
|
}
|
|
}
|
|
return n, nil
|
|
}
|
|
|
|
func (tx *Tx) OffsetRange(name string, offset, start, endx uint64) (*roaring.Bitmap, error) {
|
|
if lowbits(offset) != 0 {
|
|
vprint.PanicOn("offset must not contain low bits")
|
|
} else if lowbits(start) != 0 {
|
|
vprint.PanicOn("range start must not contain low bits")
|
|
} else if lowbits(endx) != 0 {
|
|
vprint.PanicOn("range endx must not contain low bits")
|
|
}
|
|
|
|
tx.mu.RLock()
|
|
defer tx.mu.RUnlock()
|
|
|
|
c, err := tx.cursor(name)
|
|
if err == ErrBitmapNotFound {
|
|
return roaring.NewSliceBitmap(), nil
|
|
} else if err != nil {
|
|
return nil, err
|
|
}
|
|
defer c.Close()
|
|
|
|
other := roaring.NewSliceBitmap()
|
|
off := highbits(offset)
|
|
hi0, hi1 := highbits(start), highbits(endx)
|
|
|
|
if _, err := c.Seek(hi0); err == io.EOF {
|
|
return other, nil
|
|
} else if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
for {
|
|
if err := c.Next(); err == io.EOF {
|
|
break
|
|
} else if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
elem := &c.stack.elems[c.stack.top]
|
|
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.
|
|
if ckey >= hi1 {
|
|
break
|
|
}
|
|
other.Containers.Put(off+(ckey-hi0), toContainer(cell, tx))
|
|
}
|
|
return other, nil
|
|
}
|
|
|
|
// containerFilter is like ContainerIterator, but implements ApplyFilter and
|
|
// also ApplyRewriter, depending on which is provided to it.
|
|
type containerFilter struct {
|
|
cursor *Cursor
|
|
name string
|
|
filter roaring.BitmapFilter
|
|
rewriter roaring.BitmapRewriter
|
|
tx *Tx
|
|
header roaring.Container
|
|
body [8192]byte
|
|
}
|
|
|
|
func (s *containerFilter) Close() {
|
|
// note that the cursor gets put back in the pool, but that cursor.Close
|
|
// zeroes out the cursor's tx for us.
|
|
s.cursor.Close()
|
|
s.cursor = nil
|
|
s.tx = nil
|
|
s.filter = nil
|
|
s.rewriter = nil
|
|
containerFilterPool.Put(s)
|
|
}
|
|
|
|
func (s *containerFilter) ApplyFilter() (err error) {
|
|
var minKey roaring.FilterKey
|
|
var cell leafCell
|
|
if s.filter == nil {
|
|
return errors.New("can't apply filter without a filter")
|
|
}
|
|
for err := s.cursor.Next(); err == nil; err = s.cursor.Next() {
|
|
elem := &s.cursor.stack.elems[s.cursor.stack.top]
|
|
leafPage, _, err := s.cursor.tx.readPage(elem.pgno)
|
|
if err != nil {
|
|
return fmt.Errorf("reading from pgno %d applying filter: %s", elem.pgno, err)
|
|
}
|
|
readLeafCellInto(&cell, leafPage, elem.index)
|
|
key := roaring.FilterKey(cell.Key)
|
|
if key < minKey {
|
|
continue
|
|
}
|
|
res := s.filter.ConsiderKey(key, int32(cell.BitN))
|
|
if res.Err != nil {
|
|
return res.Err
|
|
}
|
|
if res.YesKey <= key && res.NoKey <= key {
|
|
data := intoContainer(cell, s.cursor.tx, &s.header, s.body[:])
|
|
res = s.filter.ConsiderData(key, data)
|
|
if res.Err != nil {
|
|
return res.Err
|
|
}
|
|
}
|
|
minKey = res.NoKey
|
|
if minKey > key+1 {
|
|
_, err := s.cursor.Seek(uint64(minKey))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *containerFilter) ApplyRewriter() (err error) {
|
|
var minKey roaring.FilterKey
|
|
var cell leafCell
|
|
if s.rewriter == nil {
|
|
return errors.New("can't apply rewriter without a rewriter")
|
|
}
|
|
var dirty bool
|
|
var key roaring.FilterKey
|
|
var writeback roaring.ContainerWriteback = func(updateKey roaring.FilterKey, data *roaring.Container) (err error) {
|
|
dirty = true
|
|
var exact bool
|
|
if updateKey != key {
|
|
exact, err = s.cursor.Seek(uint64(updateKey))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
} else {
|
|
exact = true
|
|
}
|
|
if data.N() == 0 {
|
|
if exact {
|
|
err = s.cursor.deleteLeafCell(uint64(updateKey))
|
|
key = ^roaring.FilterKey(0)
|
|
}
|
|
// if we don't delete, we aren't changing our situation at all
|
|
} else {
|
|
cell = ConvertToLeafArgs(uint64(updateKey), data)
|
|
err = s.cursor.putLeafCell(cell)
|
|
key = ^roaring.FilterKey(0)
|
|
}
|
|
return err
|
|
}
|
|
for err := s.cursor.Next(); err == nil; err = s.cursor.Next() {
|
|
elem := &s.cursor.stack.elems[s.cursor.stack.top]
|
|
leafPage, _, err := s.cursor.tx.readPage(elem.pgno)
|
|
if err != nil {
|
|
return fmt.Errorf("reading from pgno %d applying rewriter: %s", elem.pgno, err)
|
|
}
|
|
readLeafCellInto(&cell, leafPage, elem.index)
|
|
key = roaring.FilterKey(cell.Key)
|
|
if key < minKey {
|
|
continue
|
|
}
|
|
res := s.rewriter.ConsiderKey(key, int32(cell.BitN))
|
|
if res.Err != nil {
|
|
return res.Err
|
|
}
|
|
if res.YesKey <= key && res.NoKey <= key {
|
|
data, err := intoWritableContainer(cell, s.cursor.tx, &s.header, s.body[:])
|
|
if err != nil {
|
|
return fmt.Errorf("applying rewriter: %s", err)
|
|
}
|
|
res = s.rewriter.RewriteData(key, data, writeback)
|
|
if res.Err != nil {
|
|
return res.Err
|
|
}
|
|
}
|
|
minKey = res.NoKey
|
|
// if the callback did any writing, we need to reset our cursor,
|
|
// and if the next key is far away, we should also reset our cursor.
|
|
//
|
|
// In practice the "key+64" probably comes out to "we've been told
|
|
// we're done".
|
|
if dirty || minKey > (key+64) {
|
|
dirty = false
|
|
_, err := s.cursor.Seek(uint64(minKey))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
// notify rewriter that we're done, telling them the last key we
|
|
// processed.
|
|
res := s.rewriter.RewriteData(^roaring.FilterKey(0), nil, writeback)
|
|
return res.Err
|
|
}
|
|
|
|
// containerIterator wraps Cursor to implement roaring.ContainerIterator.
|
|
type containerIterator struct {
|
|
cursor *Cursor
|
|
}
|
|
|
|
// Close must be called when the client is done
|
|
// with the containerIterator so that the internal
|
|
// Cursor can be recycled.
|
|
func (itr *containerIterator) Close() {
|
|
itr.cursor.Close()
|
|
}
|
|
|
|
// Next moves the iterator to the next container.
|
|
func (itr *containerIterator) Next() bool {
|
|
err := itr.cursor.Next()
|
|
return err == nil
|
|
}
|
|
|
|
// Value returns the current key & container.
|
|
func (itr *containerIterator) Value() (uint64, *roaring.Container) {
|
|
elem := &itr.cursor.stack.elems[itr.cursor.stack.top]
|
|
leafPage, _, _ := itr.cursor.tx.readPage(elem.pgno)
|
|
cell := readLeafCell(leafPage, elem.index)
|
|
return cell.Key, toContainer(cell, itr.cursor.tx)
|
|
}
|
|
|
|
// always returns false for Next()
|
|
type emptyContainerIterator struct{}
|
|
|
|
func (si *emptyContainerIterator) Close() {}
|
|
|
|
func (si *emptyContainerIterator) Next() bool {
|
|
return false
|
|
}
|
|
|
|
func (si *emptyContainerIterator) Value() (uint64, *roaring.Container) {
|
|
vprint.PanicOn("emptyContainerIterator never has any Values")
|
|
return 0, nil
|
|
}
|
|
|
|
func (tx *Tx) ImportRoaringBits(name string, itr roaring.RoaringIterator, clear bool, log bool, rowSize uint64) (changed int, rowSet map[uint64]int, err error) {
|
|
// begin write boilerplate
|
|
if tx.db == nil {
|
|
err = ErrTxClosed
|
|
return
|
|
} else if !tx.writable {
|
|
err = ErrTxNotWritable
|
|
return
|
|
} else if name == "" {
|
|
err = ErrBitmapNameRequired
|
|
return
|
|
}
|
|
|
|
tx.mu.Lock()
|
|
defer tx.mu.Unlock()
|
|
|
|
if err = tx.createBitmapIfNotExists(name); err != nil {
|
|
return
|
|
}
|
|
// end write boilerplate
|
|
|
|
n := itr.Len()
|
|
if n == 0 {
|
|
return
|
|
}
|
|
rowSet = make(map[uint64]int)
|
|
|
|
var currRow uint64
|
|
|
|
cur, err := tx.cursor(name)
|
|
if err != nil {
|
|
return changed, rowSet, err
|
|
}
|
|
defer cur.Close()
|
|
|
|
for itrKey, synthC := itr.NextContainer(); synthC != nil; itrKey, synthC = itr.NextContainer() {
|
|
|
|
if rowSize != 0 {
|
|
currRow = itrKey / rowSize
|
|
}
|
|
nsynth := int(synthC.N())
|
|
if nsynth == 0 {
|
|
continue
|
|
}
|
|
// INVAR: nsynth > 0
|
|
|
|
// Find existing container, if any.
|
|
var oldC *roaring.Container
|
|
if exact, err := cur.Seek(itrKey); err != nil {
|
|
return changed, rowSet, err
|
|
} else if exact {
|
|
elem := &cur.stack.elems[cur.stack.top]
|
|
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 {
|
|
// no container at the itrKey in badger (or all zero container).
|
|
if clear {
|
|
// changed of 0 and empty rowSet is perfect, no need to change the defaults.
|
|
continue
|
|
} else {
|
|
changed += nsynth
|
|
rowSet[currRow] += nsynth
|
|
|
|
if err := tx.putContainerWithCursor(cur, itrKey, synthC); err != nil {
|
|
return changed, rowSet, err
|
|
}
|
|
continue
|
|
}
|
|
}
|
|
|
|
if clear {
|
|
|
|
existN := oldC.N() // number of bits set in the old container
|
|
newC := oldC.Difference(synthC)
|
|
|
|
// update rowSet and changes
|
|
if newC.N() == existN {
|
|
// INVAR: do changed need adjusting? nope. same bit count,
|
|
// so no change could have happened.
|
|
continue
|
|
} else {
|
|
changes := int(existN - newC.N())
|
|
changed += changes
|
|
rowSet[currRow] -= changes
|
|
err = tx.putContainerWithCursor(cur, itrKey, newC)
|
|
if err != nil {
|
|
return
|
|
}
|
|
continue
|
|
}
|
|
} else {
|
|
// setting bits
|
|
|
|
existN := oldC.N()
|
|
if existN == roaring.MaxContainerVal+1 {
|
|
// completely full container already, set will do nothing. so changed of 0 default is perfect.
|
|
continue
|
|
}
|
|
if existN == 0 {
|
|
// can nsynth be zero? No, because of the continue/invariant above where nsynth > 0
|
|
changed += nsynth
|
|
rowSet[currRow] += nsynth
|
|
err = tx.putContainerWithCursor(cur, itrKey, synthC)
|
|
if err != nil {
|
|
return
|
|
}
|
|
continue
|
|
}
|
|
|
|
newC := roaring.Union(oldC, synthC) // UnionInPlace was giving us crashes on overly large containers.
|
|
|
|
if roaring.ContainerType(newC) == roaring.ContainerBitmap {
|
|
newC.Repair() // update the bit-count so .n is valid. b/c UnionInPlace doesn't update it.
|
|
}
|
|
if newC.N() != existN {
|
|
changes := int(newC.N() - existN)
|
|
changed += changes
|
|
rowSet[currRow] += changes
|
|
|
|
err = tx.putContainerWithCursor(cur, itrKey, newC)
|
|
if err != nil {
|
|
vprint.PanicOn(err)
|
|
return
|
|
}
|
|
continue
|
|
}
|
|
}
|
|
}
|
|
return
|
|
}
|
|
|
|
// flush writes the dirty pages & meta page to the WAL.
|
|
func (tx *Tx) flush() error {
|
|
w := bufio.NewWriterSize(tx.db.walFile, 65536)
|
|
|
|
// Write non-bitmap pages to WAL.
|
|
for _, pgno := range dirtyPageMapKeys(tx.dirtyPages) {
|
|
walID, err := tx.writeToWAL(w, tx.dirtyPages[pgno])
|
|
if err != nil {
|
|
return fmt.Errorf("write page to wal: %w", err)
|
|
}
|
|
tx.pageMap = tx.pageMap.Set(pgno, walID)
|
|
}
|
|
|
|
// Write bitmap headers & pages to WAL.
|
|
//
|
|
// We need to write a bitmap header before each such page. We only allocate
|
|
// one header, and we reuse it, because each write is flushing it out to
|
|
// disk, and it doesn't get stored in-memory.
|
|
var hdr []byte
|
|
if len(tx.dirtyBitmapPages) > 0 {
|
|
hdr = allocPage()
|
|
}
|
|
for _, pgno := range dirtyPageMapKeys(tx.dirtyBitmapPages) {
|
|
// Write header page.
|
|
writePageNo(hdr[:], pgno)
|
|
writeFlags(hdr[:], PageTypeBitmapHeader)
|
|
if _, err := tx.writeToWAL(w, hdr); err != nil {
|
|
return fmt.Errorf("write bitmap header page to wal: %w", err)
|
|
}
|
|
|
|
// Write bitmap page.
|
|
walID, err := tx.writeToWAL(w, tx.dirtyBitmapPages[pgno])
|
|
if err != nil {
|
|
return fmt.Errorf("write bitmap page to wal: %w", err)
|
|
}
|
|
tx.pageMap = tx.pageMap.Set(pgno, walID)
|
|
}
|
|
|
|
// At this point, it is safe to nil out the dirtyPages and
|
|
// dirtyBitmapPages objects. We don't. The reason we don't is that
|
|
// we should never have a Tx lasting for long anyway -- even if we
|
|
// end up holding the write lock for a checkpoint, we don't keep the
|
|
// associated Tx around. If we nil those out, then a few stray Tx
|
|
// objects sticking around won't stick out in a heap profile. If we
|
|
// leave them alone, they'll stick out in a heap profile. I think on
|
|
// the whole that's better for further observability and debugging.
|
|
|
|
// Write meta page to WAL.
|
|
walID, err := tx.writeToWAL(w, tx.meta[:])
|
|
if err != nil {
|
|
return fmt.Errorf("write meta page to wal: %w", err)
|
|
}
|
|
tx.pageMap = tx.pageMap.Set(uint32(0), walID)
|
|
|
|
// Flush & sync WAL.
|
|
if err := w.Flush(); err != nil {
|
|
return fmt.Errorf("flush wal: %w", err)
|
|
} else if err := tx.db.fsyncWAL(tx.db.walFile); err != nil {
|
|
return fmt.Errorf("sync wal: %w", err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (tx *Tx) writeToWAL(w io.Writer, page []byte) (walID int64, err error) {
|
|
// Determine next WAL ID from cached meta page.
|
|
walID = readMetaWALID(tx.meta[:]) + 1
|
|
|
|
// Update WAL ID on cached meta page.
|
|
writeMetaWALID(tx.meta[:], walID)
|
|
|
|
// Append to WAL and increment WAL size.
|
|
if _, err := w.Write(page); err != nil {
|
|
return 0, err
|
|
}
|
|
tx.walPageN++
|
|
|
|
return walID, nil
|
|
}
|
|
|
|
// Pages returns meta & record data for a list of pages.
|
|
func (tx *Tx) Pages(pgnos []uint32) ([]Page, error) {
|
|
// Read page info for all pages in the database.
|
|
infos, err := tx.PageInfos()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// Loop over each requested page number and extract additional data.
|
|
var pages []Page
|
|
for _, pgno := range pgnos {
|
|
buf, _, err := tx.readPage(pgno)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
switch info := infos[pgno].(type) {
|
|
case *MetaPageInfo:
|
|
pages = append(pages, &MetaPage{MetaPageInfo: info})
|
|
|
|
case *RootRecordPageInfo:
|
|
records, err := readRootRecords(buf)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
pages = append(pages, &RootRecordPage{RootRecordPageInfo: info, Records: records})
|
|
|
|
case *LeafPageInfo:
|
|
page := &LeafPage{LeafPageInfo: info}
|
|
cells := make([]leafCell, page.CellN)
|
|
for _, cell := range readLeafCells(buf, cells) {
|
|
other := &LeafCell{
|
|
Key: cell.Key,
|
|
Type: cell.Type,
|
|
}
|
|
|
|
switch cell.Type {
|
|
case ContainerTypeArray, ContainerTypeRLE:
|
|
other.Values = cell.Values(tx)
|
|
case ContainerTypeBitmapPtr:
|
|
other.Pgno = toPgno(cell.Data)
|
|
}
|
|
|
|
page.Cells = append(page.Cells, other)
|
|
}
|
|
pages = append(pages, page)
|
|
|
|
case *BranchPageInfo:
|
|
page := &BranchPage{BranchPageInfo: info}
|
|
for _, cell := range readBranchCells(buf) {
|
|
page.Cells = append(page.Cells, &BranchCell{
|
|
Key: cell.LeftKey,
|
|
Flags: cell.Flags,
|
|
Pgno: cell.ChildPgno,
|
|
})
|
|
}
|
|
pages = append(pages, page)
|
|
|
|
case *BitmapPageInfo:
|
|
pages = append(pages, &BitmapPage{
|
|
BitmapPageInfo: info,
|
|
Values: bitmapValues(toArray64(buf)),
|
|
})
|
|
|
|
case *FreePageInfo:
|
|
pages = append(pages, &FreePage{FreePageInfo: info})
|
|
|
|
default:
|
|
vprint.PanicOn(fmt.Sprintf("invalid page info type %T", info))
|
|
}
|
|
}
|
|
|
|
return pages, nil
|
|
}
|
|
|
|
// PageInfos returns meta data about all pages in the database.
|
|
func (tx *Tx) PageInfos() ([]PageInfo, error) {
|
|
var errorList ErrorList
|
|
|
|
infos := make([]PageInfo, tx.PageN())
|
|
|
|
// Read meta page info.
|
|
metaInfo, err := tx.metaPageInfo()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
infos[0] = metaInfo
|
|
|
|
// Traverse root record linked list.
|
|
for pgno := metaInfo.RootRecordPageNo; pgno != 0; {
|
|
info, err := tx.rootRecordPageInfo(pgno)
|
|
if err != nil {
|
|
errorList.Append(err)
|
|
break
|
|
}
|
|
infos[pgno] = info
|
|
pgno = info.Next
|
|
}
|
|
|
|
// Traverse freelist and mark pages as in-use.
|
|
if err := tx.walkPageInfo(infos, metaInfo.FreelistPageNo, "freelist"); err != nil {
|
|
errorList.Append(err)
|
|
}
|
|
|
|
// Traverse every b-tree and mark pages as in-use.
|
|
records, err := tx.RootRecords()
|
|
if err != nil {
|
|
errorList.Append(err)
|
|
} else {
|
|
for itr := records.Iterator(); !itr.Done(); {
|
|
name, pgno, _ := itr.Next()
|
|
|
|
if err := tx.walkPageInfo(infos, pgno, name); err != nil {
|
|
errorList.Append(err)
|
|
}
|
|
}
|
|
}
|
|
|
|
// Build page info objects for each free page.
|
|
freePageSet, err := tx.freePageSet()
|
|
if err != nil {
|
|
errorList.Append(err)
|
|
} else {
|
|
for pgno := range freePageSet {
|
|
infos[pgno] = &FreePageInfo{Pgno: pgno}
|
|
}
|
|
}
|
|
|
|
return infos, errorList.Err()
|
|
}
|
|
|
|
// metaPageInfo returns page metadata for the meta page.
|
|
func (tx *Tx) metaPageInfo() (*MetaPageInfo, error) {
|
|
buf, _, err := tx.readPage(0)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return &MetaPageInfo{
|
|
Pgno: 0,
|
|
Magic: readMetaMagic(buf),
|
|
PageN: readMetaPageN(buf),
|
|
WALID: readMetaWALID(buf),
|
|
RootRecordPageNo: readMetaRootRecordPageNo(buf),
|
|
FreelistPageNo: readMetaFreelistPageNo(buf),
|
|
}, nil
|
|
}
|
|
|
|
// rootRecordPageInfo returns page metadata for a root record page.
|
|
func (tx *Tx) rootRecordPageInfo(pgno uint32) (*RootRecordPageInfo, error) {
|
|
buf, _, err := tx.readPage(pgno)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return &RootRecordPageInfo{
|
|
Pgno: pgno,
|
|
Next: WalkRootRecordPages(buf),
|
|
}, nil
|
|
}
|
|
|
|
func (tx *Tx) walkPageInfo(infos []PageInfo, root uint32, name string) error {
|
|
var errorList ErrorList
|
|
|
|
if err := tx.walkTree(root, 0, func(pgno, parent, typ uint32, err error) error {
|
|
if err != nil {
|
|
errorList.Append(err)
|
|
return nil
|
|
}
|
|
|
|
buf, _, err := tx.readPage(pgno)
|
|
if err != nil {
|
|
errorList.Append(fmt.Errorf("cannot read page: pgno=%d parent=%d typ=%d err=%d", pgno, parent, typ, err))
|
|
return nil
|
|
}
|
|
|
|
switch typ {
|
|
case PageTypeLeaf:
|
|
infos[pgno] = &LeafPageInfo{
|
|
Pgno: pgno,
|
|
Parent: parent,
|
|
Tree: name,
|
|
Flags: readFlags(buf),
|
|
CellN: readCellN(buf),
|
|
}
|
|
case PageTypeBranch:
|
|
infos[pgno] = &BranchPageInfo{
|
|
Pgno: pgno,
|
|
Parent: parent,
|
|
Tree: name,
|
|
Flags: readFlags(buf),
|
|
CellN: readCellN(buf),
|
|
}
|
|
case PageTypeBitmap:
|
|
infos[pgno] = &BitmapPageInfo{
|
|
Pgno: pgno,
|
|
Parent: parent,
|
|
Tree: name,
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}); err != nil {
|
|
errorList.Append(err)
|
|
}
|
|
|
|
return errorList.Err()
|
|
}
|
|
|
|
// PageData returns the raw page data for a single page.
|
|
func (tx *Tx) PageData(pgno uint32) ([]byte, error) {
|
|
buf, _, err := tx.readPage(pgno)
|
|
return buf, err
|
|
}
|
|
|
|
func (tx *Tx) GetSortedFieldViewList() (fvs []txkey.FieldView, _ error) {
|
|
records, err := tx.RootRecords()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
it := records.Iterator()
|
|
for !it.Done() {
|
|
k, _, _ := it.Next()
|
|
root := k
|
|
fv := txkey.FieldViewFromPrefix([]byte(root))
|
|
fvs = append(fvs, fv)
|
|
}
|
|
return
|
|
}
|
|
|
|
func (tx *Tx) DebugInfo() *TxDebugInfo {
|
|
return &TxDebugInfo{
|
|
Ptr: fmt.Sprintf("%p", tx),
|
|
Writable: tx.writable,
|
|
Stack: string(tx.stack),
|
|
}
|
|
}
|
|
|
|
type TxDebugInfo struct {
|
|
Ptr string `json:"ptr"`
|
|
Writable bool `json:"writable"`
|
|
Stack string `json:"stack,omitempty"`
|
|
}
|
|
|
|
// SnapshotReader returns a reader that provides a snapshot for the current database state.
|
|
func (tx *Tx) SnapshotReader() (io.Reader, error) {
|
|
if tx.db == nil {
|
|
return nil, ErrTxClosed
|
|
}
|
|
return &snapshotReader{tx: tx}, nil
|
|
}
|
|
|
|
type snapshotReader struct {
|
|
tx *Tx
|
|
pgno uint32
|
|
}
|
|
|
|
func (r *snapshotReader) Read(p []byte) (n int, err error) {
|
|
// Exit if we are past the end of the database.
|
|
if r.pgno >= readMetaPageN(r.tx.meta[:]) {
|
|
return 0, io.EOF
|
|
}
|
|
|
|
// Otherwise look up the page data from mmap or page cache and copy it out.
|
|
buf, _, err := r.tx.readPage(r.pgno)
|
|
if err != nil {
|
|
return 0, err
|
|
} else if len(p) < len(buf) {
|
|
return 0, io.ErrShortBuffer
|
|
}
|
|
copy(p, buf)
|
|
|
|
// Increment the page number.
|
|
r.pgno++
|
|
|
|
return len(buf), nil
|
|
}
|
|
|
|
type PageInfo interface {
|
|
pageInfo()
|
|
}
|
|
|
|
func (*MetaPageInfo) pageInfo() {}
|
|
func (*RootRecordPageInfo) pageInfo() {}
|
|
func (*LeafPageInfo) pageInfo() {}
|
|
func (*BranchPageInfo) pageInfo() {}
|
|
func (*BitmapPageInfo) pageInfo() {}
|
|
func (*FreePageInfo) pageInfo() {}
|
|
|
|
type MetaPageInfo struct {
|
|
Pgno uint32
|
|
Magic []byte
|
|
PageN uint32
|
|
WALID int64
|
|
RootRecordPageNo uint32
|
|
FreelistPageNo uint32
|
|
}
|
|
|
|
type RootRecordPageInfo struct {
|
|
Pgno uint32
|
|
Next uint32
|
|
}
|
|
|
|
type LeafPageInfo struct {
|
|
Pgno uint32
|
|
Parent uint32
|
|
Tree string
|
|
Flags uint32
|
|
CellN int
|
|
}
|
|
|
|
type BranchPageInfo struct {
|
|
Pgno uint32
|
|
Parent uint32
|
|
Tree string
|
|
Flags uint32
|
|
CellN int
|
|
}
|
|
|
|
type BitmapPageInfo struct {
|
|
Pgno uint32
|
|
Parent uint32
|
|
Tree string
|
|
}
|
|
|
|
type FreePageInfo struct {
|
|
Pgno uint32
|
|
}
|
|
|
|
type Page interface {
|
|
page()
|
|
}
|
|
|
|
func (*MetaPage) page() {}
|
|
func (*RootRecordPage) page() {}
|
|
func (*LeafPage) page() {}
|
|
func (*BranchPage) page() {}
|
|
func (*BitmapPage) page() {}
|
|
func (*FreePage) page() {}
|
|
|
|
type MetaPage struct {
|
|
*MetaPageInfo
|
|
}
|
|
|
|
type RootRecordPage struct {
|
|
*RootRecordPageInfo
|
|
Records []*RootRecord
|
|
}
|
|
|
|
type LeafPage struct {
|
|
*LeafPageInfo
|
|
Cells []*LeafCell
|
|
}
|
|
|
|
// LeafCell represents a leaf cell in the public API.
|
|
type LeafCell struct {
|
|
Key uint64
|
|
Type ContainerType
|
|
Pgno uint32 // bitmap pointer only
|
|
Values []uint16 // array & rle containers only
|
|
}
|
|
|
|
type BranchPage struct {
|
|
*BranchPageInfo
|
|
Cells []*BranchCell
|
|
}
|
|
|
|
// BranchCell represents a branch cell in the public API.
|
|
type BranchCell struct {
|
|
Key uint64
|
|
Flags uint32
|
|
Pgno uint32
|
|
}
|
|
|
|
type BitmapPage struct {
|
|
*BitmapPageInfo
|
|
Values []uint16
|
|
}
|
|
|
|
type FreePage struct {
|
|
*FreePageInfo
|
|
}
|
|
|
|
// dirtyPageMapKeys returns a sorted slice slice of keys for a dirty page map.
|
|
func dirtyPageMapKeys(m map[uint32][]byte) []uint32 {
|
|
a := make([]uint32, 0, len(m))
|
|
for k := range m {
|
|
a = append(a, k)
|
|
}
|
|
sort.Sort(uint32Slice(a))
|
|
return a
|
|
}
|
|
|
|
type uint32Slice []uint32
|
|
|
|
func (p uint32Slice) Len() int { return len(p) }
|
|
func (p uint32Slice) Less(i, j int) bool { return p[i] < p[j] }
|
|
func (p uint32Slice) Swap(i, j int) { p[i], p[j] = p[j], p[i] }
|